LabHub
배우기 러닝패스 코스

Apache Flink — Running Streams on a Real Engine

Does the sequence continue after stop and resume?

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

순번 원천(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.sqlSTOP 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.jsonsavepoint_path·last_id_before_stop·first_id_after_resume·fresh_run_overlap·retained_checkpoint·leftover_inprogress_files 를 적으세요.

참고

체크포인트를 켠 순번 잡을 제출한다

flink-up/root/flink/checkpoint/run1.sql 에 체크포인트 간격 1 s, execution.checkpointing.dir = file:///root/flink/checkpoint/ckpt, execution.checkpointing.savepoint-dir = file:///root/flink/checkpoint/sp, 잡 이름 flk-seq-1, datagen 순번 원천, file:///root/flink/checkpoint/out1 로 가는 csv 싱크와 INSERT 를 쓰고 sql-client.sh -f run1.sql > run1.out 2>&1/root/flink/checkpoint/run1.out 을 만드세요.

SET 문은 INSERT 보다 앞에 둬야 그 잡에 먹습니다. 순번 원천은 fields.id.kind = sequence 이고 시작·끝을 fields.id.start · fields.id.end 로 줍니다. 초당 행 수를 작게(20) 두어야 파일이 몇 개 안 생겨 눈으로 따라가기 쉽습니다. 출력에 Job ID 가 보이면 잡은 뒤에서 계속 돌고 있습니다.

도는 중의 싱크 디렉터리를 본다

잡이 도는 동안 ls -A /root/flink/checkpoint/out1 의 출력을 /root/flink/checkpoint/files-running.txt 에 저장하세요. 확정된 part-… 와 쓰는 중인 .part-…inprogress… 가 함께 보여야 합니다.

체크포인트가 한 번 완료돼야 첫 part 파일이 확정됩니다. 1초 간격이면 몇 초 기다리면 됩니다. 점으로 시작하는 파일은 -A 없이는 안 보이고, 이 파드에서 재 보면 체크포인트 사이에 잠깐씩만 나타납니다 — 0.1초 간격으로 여러 번 목록을 받아 둘이 함께 보일 때 저장하세요. 채점기는 목록의 확정 파일이 지금도 out1 에 있는지 대조합니다.

체크포인트 기록을 REST 로 받는다

run1 잡의 /jobs/<jid>/checkpoints/root/flink/checkpoint/checkpoints.json, /jobs/<jid>/checkpoints/config/root/flink/checkpoint/checkpoint-config.json 에 저장하세요. 완료된 체크포인트가 하나 이상이고, 마지막 것의 경로가 이 잡의 ckpt/<jid>/chk-<번호> 여야 합니다.

잡 id 는 run1.out 의 Job ID 줄에 있습니다. checkpoints 응답의 latest.completed 에 마지막 완료 체크포인트의 번호와 external_path 가, config 응답에 모드(정확히 한 번)·간격(밀리초)·보존 여부(externalization)가 있습니다.

세이브포인트를 찍고 멈춘다

/root/flink/checkpoint/stop.sql 에 세이브포인트 디렉터리 SET 과 STOP JOB '<run1 의 jid>' WITH SAVEPOINT; 를 쓰고 출력을 /root/flink/checkpoint/stop.out 에 저장하세요. 멈춘 뒤 out1 에는 점 파일이 없고, 확정 파일의 id 가 1 부터 빈틈·중복 없이 이어져야 합니다.

STOP JOB 은 sql-client 에서 도는 문장이고 세이브포인트 경로를 표 한 칸으로 돌려줍니다. 경로는 sp/savepoint-<잡 id 앞 6자리>-… 모양입니다. 세이브포인트가 완료되는 순간 쓰던 파일도 확정되므로 멈춘 뒤에는 점 파일이 없어야 합니다 — 남아 있다면 취소로 멈춘 것입니다.

세이브포인트에서 되살려 이어 쓴다

맨 앞에 SET 'execution.state-recovery.path' = '<stop.out 의 세이브포인트>'; 를 두고 잡 이름 flk-seq-2, 싱크 file:///root/flink/checkpoint/out2 로 바꾼 /root/flink/checkpoint/run2.sql 을 제출해 /root/flink/checkpoint/run2.out 을 남기세요. 몇 초 뒤 REST(POST /jobs/<jid>/stop)로 세이브포인트를 찍으며 멈추고, COMPLETED 가 된 상태 응답을 /root/flink/checkpoint/stop2.json 에 저장하세요. out2 의 첫 id 는 out1 의 마지막 바로 다음이어야 합니다.

세이브포인트에는 원천이 다음에 낼 번호도 상태로 들어 있습니다. 그래서 싱크 경로를 바꿔도 번호는 이어집니다. REST 로 멈추면 곧바로 request-id 만 돌려주므로, /jobs//savepoints/ 를 COMPLETED 가 될 때까지 다시 받아 저장합니다.

세이브포인트 없이 새로 돌리고 취소한다

복원 경로 없이, 잡 이름 flk-seq-3 · 싱크 file:///root/flink/checkpoint/out3 · SET 'execution.checkpointing.externalized-checkpoint-retention' = 'RETAIN_ON_CANCELLATION'; 을 넣은 /root/flink/checkpoint/run3.sql 을 제출해 /root/flink/checkpoint/run3.out 을 남기세요. out3 에 확정 파일이 생긴 뒤 잡을 취소하고, 취소된 잡의 /jobs/<jid>/root/flink/checkpoint/job3.json, /jobs/<jid>/checkpoints/root/flink/checkpoint/checkpoints3.json 에 저장하세요.

복원하지 않은 잡은 원천이 1 부터 다시 셉니다 — out1 과 같은 번호가 나옵니다. 취소는 세이브포인트를 찍지 않습니다. 보존 설정이 없으면 취소와 함께 체크포인트 디렉터리가 지워지므로, 채점기는 checkpoints3.json 이 가리키는 chk 디렉터리에 _metadata 가 남아 있는지 봅니다.

남긴 체크포인트에서 이어 쓴다

checkpoints3.json 의 마지막 완료 체크포인트를 복원 경로로 두고, 잡 이름 flk-seq-4 · 싱크는 같은 file:///root/flink/checkpoint/out3/root/flink/checkpoint/run4.sql 을 제출해 /root/flink/checkpoint/run4.out 을 남기세요. 몇 초 뒤 REST 로 세이브포인트를 찍으며 멈추고 상태 응답을 /root/flink/checkpoint/stop4.json 에 저장하세요. out3 의 확정 파일은 1 부터 빈틈·중복 없이 이어져야 합니다.

체크포인트 디렉터리(chk-N)도 세이브포인트처럼 복원 경로로 쓸 수 있습니다. 새 잡은 체크포인트 시점의 번호부터 다시 쓰고, 취소 때 쓰다 만 점 파일은 결과가 아닙니다. 채점기는 점으로 시작하지 않는 파일만 모아 봅니다. 파일 이름표(uuid)가 다른 두 실행의 파일이 섞여 있어야 합니다.

보고서 — 디스크에서 다시 센다

/root/flink/checkpoint/report.jsonsavepoint_path(stop.out 의 세이브포인트), last_id_before_stop(out1 확정 파일의 마지막 id), first_id_after_resume(out2 의 첫 id), fresh_run_overlap(run3 — 1 부터 다시 돈 실행 — 이 out3 에 확정한 id 중 out1·out2 에도 있는 것의 개수), retained_checkpoint(checkpoints3.json 의 마지막 완료 경로), leftover_inprogress_files(지금 out3 에 남은 점 파일 수)를 적으세요.

모든 값은 디스크에서 다시 셀 수 있습니다. 확정 파일은 이름이 part- 로 시작하고, 한 실행의 파일은 같은 이름표(uuid)를 나눠 가집니다. run3 의 파일은 1 이 든 파일과 이름표가 같습니다. 겹침은 두 id 목록의 교집합 크기입니다(comm -12).