LabHub
배우기 러닝패스 코스

Apache Flink — Running Streams on a Real Engine

Same GROUP BY: One Batch Row per Key vs. a 229-Line Streaming Changelog

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

같은 집계를 배치와 스트리밍으로 돌려 최종 상태는 같고 스트리밍은 변경 로그를 낸다는 것을 확인한다. 연산자와 싱크에 따라 변경 종류(+I · -U · +U · -D)가 어떻게 정해지는지 실행 계획과 거절 오류로 읽는다.

왜 중요한가

스트리밍 결과는 한 번 쓰고 끝나는 값이 아니라 계속 고쳐 쓰이는 표다. 그 고침을 받는 쪽이 철회(-U)를 처리하지 못하면 숫자가 부풀고, 추가만 되는 저장소는 아예 받을 수 없다. 어떤 질의가 갱신을 내는지, 어떤 싱크가 그것을 받을 수 있는지를 계획에서 미리 읽을 줄 알아야 설계 단계에서 사고를 막는다. 이 실습의 채점기는 클러스터에 묻지 않는다 — 여러분이 저장한 sql-client 출력과 싱크 파일을 읽고, 변경 로그를 원본 CSV 에서 한 줄씩 재현해 대조한다(입력이 추가만 되고 병렬도가 1 이면 로그 순서까지 정해집니다).

단계

  1. flink-up 으로 클러스터를 띄우고, /root/flink/dynamic/batch.sql 에 배치 모드로 /opt/lab/fixtures/data/dynamic_clicks.csv 를 읽어 user_idclicks(건수)·dwell(dwell 합)을 내는 SQL 을 쓴 뒤 출력을 /root/flink/dynamic/batch.out 에 저장하세요.
  2. 같은 질의를 스트리밍 모드로 돌리는 /root/flink/dynamic/stream.sql 을 만들어 출력을 /root/flink/dynamic/stream.out 에 저장하세요.
  3. 클릭 수(clicks)별 사용자 수(users)를 내는 집계 위의 집계를 스트리밍으로 돌리는 /root/flink/dynamic/nested.sql 을 만들어 출력을 /root/flink/dynamic/nested.out 에 저장하세요.
  4. /root/flink/dynamic/explain.sqlEXPLAIN CHANGELOG_MODE 두 문장(사용자별 COUNT 집계 · dwell > 60 인 행의 user_id, url 만 고르는 필터)을 쓰고 출력을 /root/flink/dynamic/explain.out 에 저장하세요.
  5. /root/flink/dynamic/fs-sink.sql 에 filesystem 커넥터 표 user_clicks_fs(user_id STRING, clicks BIGINT) 를 만들고 사용자별 COUNT 를 INSERT INTO 하게 쓴 뒤, 거절 오류가 담긴 출력을 /root/flink/dynamic/fs-sink.out 에 저장하세요.
  6. /root/flink/dynamic/to-changelog.sql 에서 사용자별 COUNT 뷰를 TO_CHANGELOG 로 바꿔 filesystem 싱크(경로 /root/flink/dynamic/changelog, 열 change, user_id, clicks)에 쓰고 출력을 /root/flink/dynamic/to-changelog.out 에 저장하세요.
  7. /root/flink/dynamic/upsert.sql 에 print 커넥터 표 retract_sinkPRIMARY KEY (user_id) NOT ENFORCED 를 둔 blackhole 커넥터 표 upsert_sink 를 만들고, 두 싱크에 같은 집계를 넣는 EXPLAIN CHANGELOG_MODE INSERT INTO ... 를 각각 돌려 출력을 /root/flink/dynamic/upsert.out 에 저장하세요.
  8. /root/flink/dynamic/report.jsonusers·retract_messages·upsert_messages·nested_deletes 를 적으세요.

참고

배치로 한 번에 센다

flink-up 으로 클러스터를 띄우고, /root/flink/dynamic/batch.sqlSET 'execution.runtime-mode' = 'batch'; · 원본 CSV 를 읽는 CREATE TABLE · user_idclicks(COUNT(*))와 dwell(SUM(dwell))을 내는 SELECT 를 쓰고, 출력을 /root/flink/dynamic/batch.out 에 저장하세요.

filesystem 커넥터에 path 는 file:///opt/lab/fixtures/data/dynamic_clicks.csv, format 은 csv 입니다. 결과 열에는 AS clicks · AS dwell 로 별칭을 붙입니다. 배치 결과에는 op 열이 없고 사용자마다 한 줄씩 나옵니다.

같은 질의를 스트리밍으로 돌린다

1단계와 같은 질의를 SET 'execution.runtime-mode' = 'streaming'; 으로 돌리는 /root/flink/dynamic/stream.sql 을 만들고 출력을 /root/flink/dynamic/stream.out 에 저장하세요. 로그를 끝까지 적용한 최종 상태가 배치 결과와 같아야 합니다.

행이 하나 올 때마다 결과 표가 바뀝니다. 처음 보는 사용자는 +I 한 줄, 이미 있던 사용자는 옛 값을 지우는 -U 와 새 값 +U 두 줄입니다. 로그 줄 수는 사용자 수 + 2 × (클릭 수 − 사용자 수) 가 되어야 합니다. 설정은 기본 그대로 둡니다.

집계 위의 집계에서 -D 를 본다

스트리밍 모드로 SELECT clicks, COUNT(*) AS users FROM (사용자별 COUNT(*) AS clicks) GROUP BY clicks 를 돌리는 /root/flink/dynamic/nested.sql 을 만들고 출력을 /root/flink/dynamic/nested.out 에 저장하세요.

사용자가 클릭 1번 칸에서 2번 칸으로 옮길 때 안쪽 집계는 -U(1) 와 +U(2) 를 냅니다. 바깥 집계는 -U 를 받아 1번 칸의 사용자 수를 하나 줄이고, 그 칸이 0 이 되면 결과 행이 사라져야 하므로 -D 를 냅니다. 안쪽 질의의 열 이름을 clicks 로 맞추세요.

계획에서 변경 종류를 읽는다

/root/flink/dynamic/explain.sqlEXPLAIN CHANGELOG_MODE 로 시작하는 문장 두 개 — 사용자별 COUNT 집계, 그리고 dwell > 60 인 행의 user_id, url 만 고르는 필터 — 를 쓰고 출력을 /root/flink/dynamic/explain.out 에 저장하세요.

EXPLAIN 은 잡을 돌리지 않고 계획만 찍습니다. CHANGELOG_MODE 를 붙이면 Optimized Physical Plan 의 노드마다 changelogMode=[...] 가 붙습니다. I 는 추가, UB 는 갱신 전, UA 는 갱신 후, D 는 삭제입니다. 필터만 있는 질의의 노드는 무엇만 낼까요?

갱신을 파일 싱크에 넣어 거절당한다

/root/flink/dynamic/fs-sink.sql 에 filesystem 커넥터 표 user_clicks_fs(user_id STRING, clicks BIGINT)(경로는 /root/flink/dynamic/fs-out, format csv)를 만들고 INSERT INTO user_clicks_fs SELECT user_id, COUNT(*) FROM clicks GROUP BY user_id; 를 쓴 뒤, 출력을 /root/flink/dynamic/fs-sink.out 에 저장하세요. 오류로 끝나는 것이 정상입니다.

파일은 한 번 쓴 줄을 고칠 수 없습니다. 옵티마이저는 싱크가 받을 수 있는 변경 종류를 먼저 묻고, 집계가 내는 UB/UA 를 받을 수 없으면 잡을 제출하기 전에 거절합니다. 오류 문장에 어느 싱크가 어느 노드의 무엇을 못 받는지가 적혀 있습니다.

변경 종류를 열로 꺼내 파일에 쓴다

/root/flink/dynamic/to-changelog.sql 에서 사용자별 COUNT 를 뷰(user_id, clicks 두 열)로 만들고, TO_CHANGELOG(input => TABLE 뷰, op => DESCRIPTOR(change)) 의 결과를 filesystem 싱크(경로 /root/flink/dynamic/changelog, 열 change STRING, user_id STRING, clicks BIGINT, format csv)에 INSERT 하세요. 출력은 /root/flink/dynamic/to-changelog.out 에 저장합니다.

TO_CHANGELOG 는 갱신 로그의 각 행을 추가(INSERT)로 바꾸고 원래의 변경 종류를 문자열 열에 담습니다 — 그래서 추가만 하는 파일 싱크도 받을 수 있습니다. INSERT 는 비동기라 table.dml-sync 를 켜야 잡이 끝난 뒤 파일이 확정된 상태로 남습니다. 다시 돌릴 때는 싱크 디렉터리를 먼저 비우세요.

키 있는 싱크 앞에서 -U 가 사라진다

/root/flink/dynamic/upsert.sql 에 print 커넥터 표 retract_sink(user_id STRING, clicks BIGINT) 와, 같은 열에 PRIMARY KEY (user_id) NOT ENFORCED 를 둔 blackhole 커넥터 표 upsert_sink 를 만드세요. 그리고 사용자별 COUNT 를 각 싱크에 넣는 EXPLAIN CHANGELOG_MODE INSERT INTO ... 두 문장을 돌려 출력을 /root/flink/dynamic/upsert.out 에 저장하세요.

print 싱크는 받는 대로 찍기만 하므로 옛 행을 지우려면 철회(UB)가 필요합니다. 키가 선언된 싱크는 키로 찾아 덮어쓰면 되므로 UB 가 필요 없습니다. 두 계획에서 GroupAggregate 의 changelogMode 를 비교해 보세요. NOT ENFORCED 는 Flink 가 키를 검사하지 않는다는 뜻입니다.

보고서 — 리트랙트와 업서트의 메시지 수

/root/flink/dynamic/report.jsonusers(사용자 수 = stream.out 의 +I 수), retract_messages(stream.out 로그의 전체 줄 수), upsert_messages(같은 결과를 업서트 싱크로 보냈다면 보냈을 메시지 수), nested_deletes(nested.out 의 -D 수)를 정수로 적으세요.

업서트 싱크는 키로 덮어쓰므로 갱신 전 값(-U)을 받을 필요가 없습니다. 같은 로그에서 어떤 op 만 남는지 생각해 보세요. 숫자는 저장한 출력의 op 열을 세면 나옵니다 — grep 으로 줄 머리의 '| +I |' 같은 모양을 셀 수 있습니다.