Apache Spark — The answer to a slow job is in the plan and the event log
Process each arriving file exactly once and drop late data
한국어 원문으로 표시합니다.
목표
착륙 폴더에 파일이 도착할 때마다 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 와 싱크 로그가 각각 무엇을 기억하는지, 둘째 절에는 왜 이름만 바뀐 파일이 중복이 되는지, 셋째 절에는 늦은 이벤트 수와 지표가 왜 다른지를 적으세요.