데이터 파이프라인 · 이벤트 시간 워터마크와 늦게 온 자료 · 실습
어제 숫자가 오늘 바뀌었다 — 워터마크와 늦게 온 자료
목표
이벤트 시간으로 집계하는 도구 wm.py 를 만든다. 처리 시간과 이벤트 시간의 벌어짐을 분포로 재고, 워터마크를 도착 순서대로 움직이고, 창이 닫힌 뒤 온 자료를 가르고, 버리는 정책과 반영하는 정책의 결과를 견주고, 닫는 순간의 값과 그 뒤의 정정을 따로 남긴다.
왜 중요한가
이 코스의 앞 실습에서 워터마크를 한 번 다뤘다. 적재한 자료의 최대 시각을 표에 적어 다음 실행이 그 뒤만 읽게 하는 것이었고, 그 워터마크가 답하는 질문은 어디까지 읽었는가 하나였다. 여기서는 그 다음을 다룬다.
모바일 앱은 신호를 잃었다가 몇 분 뒤에, 비행기 모드였던 단말은 몇 시간 뒤에 이벤트를 보낸다. 그 이벤트가 일어난 시각은 이미 지나간 시각이고, 우리는 그 시간대의 집계를 이미 내보냈다. 워터마크는 그 자료가 더 오지 않는다고 본다는 약속이지 사실이 아니다.
그래서 정할 것이 셋이다. 허용 지연을 얼마로 잡을 것인가, 약속을 어기고 온 자료를 버릴 것인가 늦은 갱신으로 반영할 것인가, 그리고 반영했다면 어제 낸 숫자와 오늘 낸 숫자의 차이를 어떻게 설명할 것인가.
시간을 재는 실습이 아니다. 이벤트에 적힌 event_time 과 ingest_time 이라는 두 칸으로 계산하는 실습이고, 프로그램이 실제로 몇 초 걸리는지는 아무 상관이 없다.
채점기는 여러분이 적어 낸 문구를 믿지 않는다. 임시 파일에 채점기가 만든 이벤트 흐름을 차려 놓고 여러분의 도구를 실제로 돌려, 창 크기와 허용 지연을 바꿔 가며 답을 대조한다. 건수와 지연 분포는 실행마다 바뀐다.
단계
1. /root/wmark/gen_stream.py 를 만들어 실행해 /root/wmark/work/stream.jsonl 을 만드세요.
2. /root/wmark/wm.py 에 skew 를 만들어 지연 분포와 순서가 어긋난 건수를 내게 하세요.
3. watermark 를 더해 도착 순서대로 워터마크를 움직이게 하세요.
4. windows 를 더해 이벤트 시간으로 창을 나눠 집계하게 하세요.
5. late 를 더해 창이 닫힌 뒤 온 자료를 가르게 하세요.
6. agg 를 더해 버리는 정책과 반영하는 정책의 결과를 각각 내게 하세요.
7. close 를 더해 닫는 순간의 값과 그 뒤의 정정을 따로 내게 하세요.
8. /root/wmark/work/watermark_report.json 과 /root/wmark/work/watermark_report.md 를 쓰세요.
참고
- 작업은 전부
/root/wmark아래에서 합니다. 자료는/root/wmark/work/stream.jsonl입니다. - 자료 한 줄은 JSON 한 덩이이고
id·key·event_time·ingest_time·amount를 담습니다. 두 시각은 에포크 초 정수입니다. 다른 칸이 더 있어도 됩니다. - 파일에 적힌 차례가 곧 도착 차례입니다. 이벤트 시간으로 다시 정렬하면 이 실습의 모든 판정이 무너집니다.
- 실행 계약:
python3 /root/wmark/wm.py <명령> <파일> [--size=초] [--lateness=초] [--policy=drop|update]. 답은 JSON 한 덩어리로 표준출력에 냅니다. 성공하면 종료 코드 0, 파일이 없으면 3, 사용법이나 정책 이름이 틀리면 2 입니다. - 창은 고정 크기이고 겹치지 않습니다. 이벤트의 창 시작은
event_time - (event_time % size)이고, 창의 범위는 시작부터 시작에 크기를 더한 값 직전까지입니다. JSON 의 키는 창 시작을 문자열로 적은 것입니다. - 지연은
ingest_time - event_time입니다. 분위수는 가장 가까운 순위로 잡습니다 — 값을 오름차순으로 놓고ceil(건수 * p / 100)번째(1부터)를 고릅니다. 보간하지 않습니다. skew <파일>응답:{"events": 정수, "out_of_order": 정수, "min_lag": 정수, "max_lag": 정수, "p50_lag": 정수, "p95_lag": 정수}. out_of_order 는 자기보다 먼저 도착한 것들의 최대 이벤트 시간보다 이벤트 시간이 이른 건수입니다.watermark <파일> --lateness=<초>응답:{"lateness": 정수, "max_event_time": 정수, "advances": 정수, "final_watermark": 정수}. advances 는 최대 이벤트 시간을 새로 갱신한 도착 건수입니다.windows <파일> --size=<초>응답:{"size": 정수, "count": 정수, "windows": {"창시작": {"events": 정수, "amount": 정수}}}. 늦고 이르고는 따지지 않고 전부 셉니다.late <파일> --size=<초> --lateness=<초>응답:{"size": 정수, "lateness": 정수, "on_time": 정수, "late": 정수, "late_by_window": {"창시작": 정수}}. 어느 이벤트가 늦었는지는 그보다 먼저 도착한 것들로만 계산한 워터마크로 판단합니다. 워터마크가 그 창의 끝 이상이면 늦은 것입니다. 첫 도착은 비교할 앞 자료가 없으므로 늦지 않습니다.agg <파일> --size --lateness --policy=drop|update응답:{"policy": 문자열, "size": 정수, "lateness": 정수, "windows": {...}, "dropped": 정수, "restated": [창시작 문자열 오름차순]}. drop 이면 늦은 자료를 빼고 dropped 로 세며 restated 는 빈 목록입니다. update 면 늦은 자료도 넣고 dropped 는 0 이며 늦은 자료가 들어간 창을 restated 에 담습니다.close <파일> --size --lateness응답:{"size": 정수, "lateness": 정수, "sealed": {...}, "corrections": [{"window": 창시작, "delta_events": 정수, "delta_amount": 정수}], "final": {...}}. sealed 는 닫히기 전에 온 것만 담은 값이고(늦은 자료만 있는 창은 0 으로 둡니다), corrections 는 창 시작 오름차순이며, final 은 둘을 더한 값입니다.- 보고서 JSON 에는 size·lateness·events·windows·on_time·late·dropped·restated·max_lag·p95_lag·policy 를 담습니다. lateness 는 1 이상 max_lag 미만으로 잡고, 그 값에서 늦게 온 자료가 1건 이상 나와야 합니다. 나오지 않으면 허용 지연을 줄이세요.
- 보고서 MD 의 절 제목은
## 무엇을 재었나## 허용 지연을 얼마로 잡았나## 늦게 온 자료를 어떻게 했나## 어제 숫자가 바뀐 이유입니다. - 공식 문서: [Flink Generating Watermarks](https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/event-time/generating_watermarks/) · [Flink Windows](https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/operators/windows/) · [python json](https://docs.python.org/3/library/json.html)
- 흔한 실수: 자료를 이벤트 시간으로 다시 정렬하기(늦은 자료가 0건으로 나옵니다), 워터마크가 뒤로 가게 두기, 버린 건수를 세지 않기, 창 시작을 문자열로 정렬하기.
단계 8개
- 늦게 오는 이벤트를 손에 쥐기
- 얼마나 벌어지는지부터 잰다
- 워터마크를 도착 순서대로 움직이기
- 이벤트 시간으로 창 나누기
- 창이 닫힌 뒤 온 자료 가려내기
- 버릴 것인가 반영할 것인가
- 닫는 순간의 값과 정정을 따로 남기기
- 바뀐 숫자를 설명하는 한 장