Apache Spark — 느린 잡의 답은 실행 계획과 이벤트 로그에 있다 · Structured Streaming · 实验
도착하는 파일을 한 번씩만 처리하고 늦은 자료를 버린다
목표
착륙 폴더에 파일이 도착할 때마다 Structured Streaming 으로 처리해 Parquet 싱크에 쓰고, 체크포인트와 싱크 로그가 '어디까지 처리했는가' 를 어떻게 기억하는지 확인한다. 이름만 다른 재전송 파일이 중복을 만드는 것을 보고 dropDuplicates 로 막은 뒤, 워터마크가 늦게 온 자료를 버리는 모습을 창 집계로 본다.
왜 중요한가
스트리밍의 어려움은 계산이 아니라 기억이다. 잡이 죽었다 다시 뜨면 무엇을 이미 했고 무엇을 아직 안 했는지 알아야 한다. Structured Streaming 은 이것을 체크포인트 폴더에 적는다. 배치를 시작하기 전에 읽을 범위를 offsets 에 쓰고, 싱크에 결과를 다 쓴 뒤 commits 에 적는다. 파일 싱크는 자기 폴더의 _spark_metadata 에 배치마다 쓴 파일 목록을 남기고, 읽는 쪽은 그 목록에 있는 파일만 결과로 본다.
그 기억에는 한계가 있다. 파일 소스는 파일 이름으로 처리 여부를 기억하므로, 협력사가 같은 내용을 다른 이름으로 다시 보내면 새 자료로 처리한다. 내용으로 중복을 지우려면 본 적 있는 키를 상태로 들고 있어야 한다.
이벤트 시각으로 창을 집계하면 또 다른 문제가 생긴다. 창을 언제 닫을 것인가. 워터마크는 '지금까지 본 가장 늦은 시각 − 허용 지연' 이고, 그보다 오래된 자료는 늦은 자료로 버리고, 끝이 워터마크를 지난 창만 결과로 낸다(append 모드). 워터마크는 배치 사이에서만 움직이므로, 자료가 몇 배치로 나뉘어 들어오느냐가 결과를 바꾼다.
단계
1. 착륙 폴더 /root/spk/stream/in 을 만들고 /data/stream/batch-01.jsonl 을 복사하세요.
2. /root/spk/stream/stream.py(앱 spk-stream-run)로 착륙 폴더를 읽어 /root/spk/stream/out/events 에 Parquet 으로 쓰는 스트림을 체크포인트 /root/spk/stream/ckpt/events, trigger(availableNow=True) 로 한 번 돌리세요.
3. batch-02.jsonl 을 착륙 폴더에 더하고 같은 스트림을 다시 돌리세요. 새 파일만 새 배치가 되어야 합니다.
4. batch-03.jsonl 과 batch-04.jsonl 을 함께 더하고 다시 돌리세요. 두 파일이 한 배치로 처리되어야 합니다.
5. 재전송 파일 batch-06.jsonl(내용은 batch-03 과 같음)을 더하고 다시 돌린 뒤, 싱크에 두 번 들어간 event_id 수를 /root/spk/stream/out/dups.txt 에 정수로 적으세요.
6. /root/spk/stream/dedup.py(앱 spk-stream-dedup)로 착륙 폴더를 dropDuplicates(["event_id"]) 해서 /root/spk/stream/out/dedup 에 쓰는 새 스트림(체크포인트 /root/spk/stream/ckpt/dedup)을 돌리세요.
7. 두 번째 착륙 폴더 /root/spk/stream/late_in 에 batch-01~batch-05 를 cp -p 로 복사하고, /root/spk/stream/window.py(앱 spk-stream-window)로 maxFilesPerTrigger=1, 워터마크 10분, 10분 창 개수 집계를 append 모드로 /root/spk/stream/out/windows 에 쓰세요(체크포인트 /root/spk/stream/ckpt/windows). 진행 기록을 /root/spk/stream/out/progress.json 에, 늦은 자료 분석을 /root/spk/stream/out/late.json 에 쓰세요.
8. /root/spk/stream/report.md 에 ## 한 번씩만 처리하기 ## 재전송과 중복 ## 늦은 자료 세 절을 쓰세요. 둘째 절에 5단계의 중복 수를, 셋째 절에 7단계의 늦은 이벤트 수를 넣으세요.
참고
- 묶음 파일:
/data/stream/batch-01.jsonl…batch-06.jsonl(각 400줄, 칸: event_id, device, event_time, value). batch-05 에는 40분쯤 늦게 도착한 이벤트가 섞여 있고, batch-06 은 batch-03 의 재전송입니다. - 스키마를 반드시 주세요:
event_id STRING, device STRING, event_time TIMESTAMP, value INT. 스트리밍 파일 소스는 스키마 추론을 기본으로 하지 않습니다. - 파일 소스는 수정 시각 순서로 파일을 집습니다. 7단계에서
cp -p로 원본의 시각(묶음마다 1분 간격)을 지켜야 순서가 묶음 번호와 같아집니다. - 싱크 로그
out/<싱크>/_spark_metadata/<배치 번호>는 첫 줄v1뒤에 파일마다 JSON 한 줄입니다. 커밋된 결과는 이 목록에 있는 파일뿐입니다. - 흔한 실수: 체크포인트 폴더를 지우고 다시 돌려 처음부터 다시 처리하는 것, 두 스트림이 한 체크포인트를 같이 쓰는 것, 7단계를 한 배치로 처리해 워터마크가 움직이지 않는 것.
- 공식 문서: [Structured Streaming Programming Guide](https://spark.apache.org/docs/4.2.0/streaming/index.html) · [Getting Started](https://spark.apache.org/docs/4.2.0/streaming/getting-started.html) · [APIs on DataFrames and Datasets](https://spark.apache.org/docs/4.2.0/streaming/apis-on-dataframes-and-datasets.html)
8个步骤
- 착륙 폴더에 첫 묶음 내려놓기
- 첫 배치 — 체크포인트와 싱크 로그
- 새 파일만 다음 배치가 된다
- 두 파일이 한 배치로
- 이름만 다른 재전송이 중복을 만든다
- dropDuplicates — 본 적 있는 id 를 기억하기
- 워터마크 — 늦은 자료는 버리고 닫힌 창만 낸다
- 기억·중복·늦음을 숫자로 남기기