Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
届くファイルを一度ずつだけ処理し、遅れたデータを捨てる
한국어 원문으로 표시합니다.
목표
착륙 폴더에 파일이 도착할 때마다 Structured Streaming 으로 처리해 Parquet 싱크에 쓰고, 체크포인트와 싱크 로그가 '어디까지 처리했는가' 를 어떻게 기억하는지 확인한다. 이름만 다른 재전송 파일이 중복을 만드는 것을 보고 dropDuplicates 로 막은 뒤, 워터마크가 늦게 온 자료를 버리는 모습을 창 집계로 본다.
왜 중요한가
스트리밍의 어려움은 계산이 아니라 기억이다. 잡이 죽었다 다시 뜨면 무엇을 이미 했고 무엇을 아직 안 했는지 알아야 한다. Structured Streaming 은 이것을 체크포인트 폴더에 적는다. 배치를 시작하기 전에 읽을 범위를 offsets 에 쓰고, 싱크에 결과를 다 쓴 뒤 commits 에 적는다. 파일 싱크는 자기 폴더의 _spark_metadata 에 배치마다 쓴 파일 목록을 남기고, 읽는 쪽은 그 목록에 있는 파일만 결과로 본다.
그 기억에는 한계가 있다. 파일 소스는 파일 이름으로 처리 여부를 기억하므로, 협력사가 같은 내용을 다른 이름으로 다시 보내면 새 자료로 처리한다. 내용으로 중복을 지우려면 본 적 있는 키를 상태로 들고 있어야 한다.
이벤트 시각으로 창을 집계하면 또 다른 문제가 생긴다. 창을 언제 닫을 것인가. 워터마크는 '지금까지 본 가장 늦은 시각 − 허용 지연' 이고, 그보다 오래된 자료는 늦은 자료로 버리고, 끝이 워터마크를 지난 창만 결과로 낸다(append 모드). 워터마크는 배치 사이에서만 움직이므로, 자료가 몇 배치로 나뉘어 들어오느냐가 결과를 바꾼다.
단계
- 착륙 폴더 /root/spk/stream/in 을 만들고
/data/stream/batch-01.jsonl을 복사하세요. - /root/spk/stream/stream.py(앱
spk-stream-run)로 착륙 폴더를 읽어 /root/spk/stream/out/events 에 Parquet 으로 쓰는 스트림을 체크포인트 /root/spk/stream/ckpt/events,trigger(availableNow=True)로 한 번 돌리세요. batch-02.jsonl을 착륙 폴더에 더하고 같은 스트림을 다시 돌리세요. 새 파일만 새 배치가 되어야 합니다.batch-03.jsonl과batch-04.jsonl을 함께 더하고 다시 돌리세요. 두 파일이 한 배치로 처리되어야 합니다.- 재전송 파일
batch-06.jsonl(내용은 batch-03 과 같음)을 더하고 다시 돌린 뒤, 싱크에 두 번 들어간event_id수를 /root/spk/stream/out/dups.txt 에 정수로 적으세요. - /root/spk/stream/dedup.py(앱
spk-stream-dedup)로 착륙 폴더를dropDuplicates(["event_id"])해서 /root/spk/stream/out/dedup 에 쓰는 새 스트림(체크포인트 /root/spk/stream/ckpt/dedup)을 돌리세요. - 두 번째 착륙 폴더 /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 에 쓰세요. - /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 · Getting Started · APIs on DataFrames and Datasets
착륙 폴더에 첫 묶음 내려놓기
착륙 폴더 /root/spk/stream/in 을 만들고 /data/stream/batch-01.jsonl 을 그 안에 복사하세요(/root/spk/stream/in/batch-01.jsonl).
스트리밍 파일 소스는 폴더를 지켜보다가 새로 나타난 파일을 다음 배치로 집습니다. 파일은 완성된 채로 한 번에 나타나야 합니다 — 쓰는 도중의 파일을 집으면 반쪽이 처리됩니다(그래서 보통 다른 곳에 쓴 뒤 옮깁니다).
첫 배치 — 체크포인트와 싱크 로그
/root/spk/stream/stream.py 를 앱 이름 spk-stream-run 으로 만들어 /root/spk/stream/in 을 스키마(event_id STRING, device STRING, event_time TIMESTAMP, value INT)로 읽고, /root/spk/stream/out/events 에 Parquet 으로 쓰는 스트림을 checkpointLocation=/root/spk/stream/ckpt/events, trigger(availableNow=True) 로 시작해 끝날 때까지 기다리세요.
availableNow 는 '지금 있는 것을 다 처리하고 멈춘다' 입니다. 끝나면 체크포인트에 offsets/0 과 commits/0 이, 싱크 폴더의 _spark_metadata/0 에 이 배치가 쓴 파일 목록이 남습니다. 채점기는 그 목록의 파일을 읽어 batch-01 의 event_id 와 같은지 봅니다.
새 파일만 다음 배치가 된다
/data/stream/batch-02.jsonl 을 착륙 폴더에 더하고 2단계의 스트림을 그대로 다시 돌리세요. 새 배치에는 batch-02 의 이벤트만 들어가야 합니다.
체크포인트가 이미 batch-01 을 처리했다고 기억하므로 다시 읽지 않습니다. 체크포인트를 지우면 기억도 사라져 처음부터 다시 처리하고, 싱크에 같은 자료가 한 벌 더 쌓입니다. 채점기는 싱크 로그에서 batch-02 의 id 만 담은 배치를 찾습니다.
두 파일이 한 배치로
batch-03.jsonl 과 batch-04.jsonl 을 함께 착륙 폴더에 더하고 스트림을 다시 돌리세요. 두 파일의 이벤트가 한 배치로 처리되어야 합니다.
배치 하나가 파일 몇 개를 집을지는 maxFilesPerTrigger 가 정하고, 주지 않으면 그때 있는 새 파일을 전부 집습니다. 배치 경계는 지연과 워터마크에 영향을 주므로 7단계에서 다시 만납니다.
이름만 다른 재전송이 중복을 만든다
재전송 파일 /data/stream/batch-06.jsonl(내용은 batch-03 과 같다)을 착륙 폴더에 더하고 스트림을 다시 돌리세요. 그다음 싱크의 커밋된 파일(각 배치의 _spark_metadata 목록)을 읽어 두 번 이상 나온 event_id 의 개수를 /root/spk/stream/out/dups.txt 에 정수로 적으세요.
파일 소스의 기억은 파일 이름입니다. 이름이 새로우니 새 자료로 처리합니다. 싱크 폴더를 그냥 읽어도 되지만(Spark 는 _spark_metadata 를 보고 읽습니다), 디렉터리의 part 파일을 직접 훑는 도구는 실패한 배치의 찌꺼기까지 집을 수 있다는 것을 기억하세요.
dropDuplicates — 본 적 있는 id 를 기억하기
/root/spk/stream/dedup.py 를 앱 이름 spk-stream-dedup 으로 만들어 같은 착륙 폴더를 읽고 dropDuplicates(["event_id"]) 한 결과를 /root/spk/stream/out/dedup 에 쓰는 새 스트림을 checkpointLocation=/root/spk/stream/ckpt/dedup, availableNow 로 돌리세요. 결과의 event_id 는 모두 한 번씩이어야 합니다.
중복 제거는 상태 연산입니다. 본 id 를 상태 저장소에 들고 있다가 같은 id 가 오면 버립니다. 워터마크 없이 하면 상태가 끝없이 자라므로, 운영에서는 withWatermark 와 함께 쓰거나 dropDuplicatesWithinWatermark 를 씁니다. 체크포인트는 스트림마다 따로입니다.
워터마크 — 늦은 자료는 버리고 닫힌 창만 낸다
/root/spk/stream/late_in 을 만들어 batch-01~batch-05 를 cp -p 로 복사하고, /root/spk/stream/window.py 를 앱 이름 spk-stream-window 로 만들어 maxFilesPerTrigger=1 로 읽고 withWatermark("event_time", "10 minutes") 뒤 10분 창으로 개수를 세어 칼럼 start,end,count 로 /root/spk/stream/out/windows 에 append 모드로 쓰세요(checkpointLocation=/root/spk/stream/ckpt/windows, availableNow). 끝난 뒤 query.recentProgress 에서 배치마다 batch·rows·watermark·dropped(상태 연산자의 numRowsDroppedByWatermark 합)를 /root/spk/stream/out/progress.json 에 목록으로 쓰고, batch-05 가 들어올 때의 워터마크보다 이른 batch-05 이벤트 수를 /root/spk/stream/out/late.json 에 {"watermark": "yyyy-MM-dd HH:mm:ss", "late_events": 정수, "dropped_rows_metric": 정수} 로 쓰세요.
워터마크는 배치가 끝날 때 '본 것 중 가장 늦은 시각 − 10분' 으로 올라가고 다음 배치부터 쓰입니다. 그래서 파일 하나씩 배치로 넣어야 batch-05 가 들어올 때 워터마크가 이미 09:29 무렵에 와 있습니다. 늦은 이벤트는 몇십 건인데 dropped 지표가 1 인 것도 보세요 — 상태 연산자 앞에서 부분 집계가 이미 창별로 한 줄로 묶었기 때문입니다.
기억·중복·늦음을 숫자로 남기기
/root/spk/stream/report.md 에 ## 한 번씩만 처리하기 ## 재전송과 중복 ## 늦은 자료 세 절을 쓰세요. 둘째 절에 5단계의 중복 수를, 셋째 절에 7단계의 late_events 와 dropped_rows_metric 을 넣으세요.
첫 절에는 체크포인트의 offsets·commits 와 싱크 로그가 각각 무엇을 기억하는지, 둘째 절에는 왜 이름만 바뀐 파일이 중복이 되는지, 셋째 절에는 늦은 이벤트 수와 지표가 왜 다른지를 적으세요.