Apache Flink — 스트림을 엔진으로 돌린다 · 동적 테이블과 변경 로그 · 실습
같은 GROUP BY, 배치 한 줄과 스트리밍 로그 229줄
목표
같은 집계를 배치와 스트리밍으로 돌려 최종 상태는 같고 스트리밍은 변경 로그를 낸다는 것을 확인한다. 연산자와 싱크에 따라 변경 종류(+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_id 별 clicks(건수)·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.sql 에 EXPLAIN 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_sink 와 PRIMARY KEY (user_id) NOT ENFORCED 를 둔 blackhole 커넥터 표 upsert_sink 를 만들고, 두 싱크에 같은 집계를 넣는 EXPLAIN CHANGELOG_MODE INSERT INTO ... 를 각각 돌려 출력을 /root/flink/dynamic/upsert.out 에 저장하세요.
8. /root/flink/dynamic/report.json 에 users·retract_messages·upsert_messages·nested_deletes 를 적으세요.
참고
- 원본 열:
click_id BIGINT, user_id STRING, url STRING, dwell INT, ts TIMESTAMP(3)(머리글 없는 CSV, 파일 순서 = 도착 순서). - SQL 파일 실행:
sql-client.sh -f 파일.sql > 파일.out 2>&1. 오류가 난 문장에서 멈추고 출력에[ERROR]가 남습니다. - 스트리밍 결과 표는 맨 앞에
op열이 붙습니다. 끝이 있는 파일을 읽으므로 스트리밍 잡도 파일을 다 읽으면 끝납니다. - 흔한 실수: 6단계를 두 번 돌리면 싱크 디렉터리에 파일이 쌓입니다. 다시 돌릴 때는 디렉터리를 먼저 비우세요.
INSERT는 기본이 비동기 제출이라 잡이 끝나기 전에 sql-client 가 나옵니다 —SET 'table.dml-sync' = 'true';를 주면 끝날 때까지 기다립니다. - 흔한 실수:
SET 'table.exec.mini-batch.enabled'같은 설정을 켜면 중간 로그가 합쳐져 줄 수가 달라집니다. 이 실습은 기본 설정 그대로 둡니다. - 공식 문서: [Dynamic Tables](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/sql-table-concepts/dynamic_tables/) · [EXPLAIN](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/utility/explain/) · [Changelog Conversion](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/queries/changelog/) · [FileSystem](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/connectors/table/filesystem/) · [Print](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/connectors/table/print/) · [BlackHole](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/connectors/table/blackhole/)
8단계
- 배치로 한 번에 센다
- 같은 질의를 스트리밍으로 돌린다
- 집계 위의 집계에서 -D 를 본다
- 계획에서 변경 종류를 읽는다
- 갱신을 파일 싱크에 넣어 거절당한다
- 변경 종류를 열로 꺼내 파일에 쓴다
- 키 있는 싱크 앞에서 -U 가 사라진다
- 보고서 — 리트랙트와 업서트의 메시지 수