Apache Flink — 스트림을 엔진으로 돌린다 · 체크포인트와 세이브포인트 · 실습
멈추고 되살려도 순번이 이어지는가
목표
순번 원천(1, 2, 3 …)을 파일 싱크에 이어 두고 체크포인트·세이브포인트로 멈췄다 되살린 뒤, 확정된 파일의 id 가 빈틈·중복 없이 이어지는지를 디스크에서 직접 확인한다. 세이브포인트 없이 다시 돌리면 무엇이 겹치는지, 취소해도 남긴 체크포인트가 무엇을 지켜 주는지도 센다.
왜 중요한가
재배포·장애 복구 때 원천의 위치, 연산자 상태, 싱크의 미확정 파일이 한 순간으로 맞지 않으면 결과가 빠지거나 두 번 나간다. 체크포인트는 이 셋을 배리어 하나로 맞추고, 파일 싱크는 체크포인트가 완료돼야 파일을 확정한다. 이 실습의 채점기는 클러스터에 묻지 않는다. 멈춘 시점의 행 수는 실행마다 달라서 정답 숫자가 없다 — 대신 확정 파일의 id 가 1 부터 빈틈·중복 없이 이어지는가, 되살린 잡이 바로 다음 id 에서 시작하는가, 세이브포인트·체크포인트 디렉터리에 _metadata 가 남았는가를 봅니다.
단계
1. flink-up 으로 클러스터를 띄우고, 체크포인트 1초 간격 · 체크포인트 디렉터리 file:///root/flink/checkpoint/ckpt · 세이브포인트 디렉터리 file:///root/flink/checkpoint/sp · 잡 이름 flk-seq-1 로 datagen 순번(id 1..100000, 초당 20행)을 file:///root/flink/checkpoint/out1 에 csv 로 쓰는 /root/flink/checkpoint/run1.sql 을 제출하고 출력을 /root/flink/checkpoint/run1.out 에 저장하세요.
2. 잡이 도는 동안 ls -A /root/flink/checkpoint/out1 결과를 /root/flink/checkpoint/files-running.txt 에 저장하세요. 확정 파일(part-…)과 쓰는 중인 파일(.part-…inprogress…)이 함께 보여야 합니다.
3. 그 잡의 /jobs/<jid>/checkpoints 를 /root/flink/checkpoint/checkpoints.json, /jobs/<jid>/checkpoints/config 를 /root/flink/checkpoint/checkpoint-config.json 에 저장하세요.
4. /root/flink/checkpoint/stop.sql 로 STOP JOB '<jid>' WITH SAVEPOINT 를 실행하고 출력을 /root/flink/checkpoint/stop.out 에 저장하세요. 멈춘 뒤 out1 에 점 파일이 남지 않아야 합니다.
5. 세이브포인트에서 되살리는 /root/flink/checkpoint/run2.sql(잡 이름 flk-seq-2, 싱크 file:///root/flink/checkpoint/out2)을 제출해 출력을 /root/flink/checkpoint/run2.out 에 두고, 몇 초 뒤 REST 로 세이브포인트를 찍으며 멈춰 완료된 상태 응답을 /root/flink/checkpoint/stop2.json 에 저장하세요.
6. 복원 없이 처음부터 도는 /root/flink/checkpoint/run3.sql(잡 이름 flk-seq-3, 싱크 file:///root/flink/checkpoint/out3, 체크포인트 보존 RETAIN_ON_CANCELLATION)을 제출해 /root/flink/checkpoint/run3.out 을 남기고, 파일이 확정된 뒤 잡을 취소하세요. 취소된 잡의 /jobs/<jid> 를 /root/flink/checkpoint/job3.json, /jobs/<jid>/checkpoints 를 /root/flink/checkpoint/checkpoints3.json 에 저장하세요.
7. 남은 체크포인트에서 되살려 같은 out3 에 이어 쓰는 /root/flink/checkpoint/run4.sql(잡 이름 flk-seq-4)을 제출해 /root/flink/checkpoint/run4.out 을 남기고, 몇 초 뒤 REST 로 세이브포인트를 찍으며 멈춰 상태 응답을 /root/flink/checkpoint/stop4.json 에 저장하세요.
8. /root/flink/checkpoint/report.json 에 savepoint_path·last_id_before_stop·first_id_after_resume·fresh_run_overlap·retained_checkpoint·leftover_inprogress_files 를 적으세요.
참고
- 원천 모양:
CREATE TABLE seq (id BIGINT) WITH ('connector' = 'datagen', 'fields.id.kind' = 'sequence', 'fields.id.start' = '1', 'fields.id.end' = '100000', 'rows-per-second' = '20'). 끝(end)을 크게 잡을수록 원천의 상태가 커집니다 — 1억으로 두면 이 파드에서는 잡이 뜨지 못하고 재시작만 되풀이합니다. - 싱크 모양:
CREATE TABLE sink (id BIGINT) WITH ('connector' = 'filesystem', 'path' = 'file:///…', 'format' = 'csv'). INSERT 는 제출만 하고 돌아오며 출력에Job ID:가 찍힙니다. - REST 로 세이브포인트와 함께 멈추기:
curl -s -X POST localhost:8081/jobs/<jid>/stop -H 'Content-Type: application/json' -d '{"drain": false, "targetDirectory": "file:///root/flink/checkpoint/sp"}'→request-id를 받아/jobs/<jid>/savepoints/<request-id>가COMPLETED가 될 때까지 다시 받습니다. - 세이브포인트 없이 취소:
curl -s -X PATCH "localhost:8081/jobs/<jid>?mode=cancel". - 되살리기: SQL 파일 맨 앞에
SET 'execution.state-recovery.path' = '<세이브포인트나 chk-N 경로>';. - 흔한 실수:
ls만 쓰면 점 파일이 안 보입니다(ls -A). 멈추지 않고 찍은 세이브포인트는 파일을 확정해 주지 않습니다. - 공식 문서: [Checkpoints](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/state/checkpoints/) · [Savepoints](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/state/savepoints/) · [JOB Statements](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/utility/job/) · [Stateful Stream Processing](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/stateful-stream-processing/) · [FileSystem 커넥터](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/connectors/table/filesystem/)
8단계
- 체크포인트를 켠 순번 잡을 제출한다
- 도는 중의 싱크 디렉터리를 본다
- 체크포인트 기록을 REST 로 받는다
- 세이브포인트를 찍고 멈춘다
- 세이브포인트에서 되살려 이어 쓴다
- 세이브포인트 없이 새로 돌리고 취소한다
- 남긴 체크포인트에서 이어 쓴다
- 보고서 — 디스크에서 다시 센다