TT Lab
Get started
Learn Learning paths Courses

Apache Flink — Running Streams on a Real Engine

Does the sequence continue after stop and resume?

Continue in TT Lab

Goal

Connect a sequence source (1, 2, 3 …) to a file sink, stop and restore it with checkpoints and savepoints, and then confirm directly from the disk whether the ids of the finalized files continue without gaps or duplicates. Also count what overlaps when you run again without a savepoint, and what the checkpoint you kept protects even after cancellation.

Why it matters

When the source position, the operator state, and the sink's unfinalized files do not line up at one moment during a redeploy or failure recovery, results go missing or go out twice. A checkpoint aligns these three with a single barrier, and a file sink finalizes files only when the checkpoint completes. The grader of this lab does not ask the cluster anything. The number of rows at the moment of stopping differs from run to run, so there is no correct number — instead it looks at whether the ids of the finalized files continue from 1 without gaps or duplicates, whether the restored job starts at the very next id, and whether _metadata remains in the savepoint and checkpoint directories.

Steps

  1. Start the cluster with flink-up and submit /root/flink/checkpoint/run1.sql, which writes a datagen sequence (id 1..100000, 20 rows per second) as csv to file:///root/flink/checkpoint/out1 with a 1-second checkpoint interval, the checkpoint directory file:///root/flink/checkpoint/ckpt, the savepoint directory file:///root/flink/checkpoint/sp, and the job name flk-seq-1, and save the output to /root/flink/checkpoint/run1.out.
  2. While the job is running, save the result of ls -A /root/flink/checkpoint/out1 to /root/flink/checkpoint/files-running.txt. Finalized files (part-…) and files being written (.part-…inprogress…) must show up together.
  3. Save that job's /jobs/<jid>/checkpoints to /root/flink/checkpoint/checkpoints.json and /jobs/<jid>/checkpoints/config to /root/flink/checkpoint/checkpoint-config.json.
  4. Run STOP JOB '<jid>' WITH SAVEPOINT with /root/flink/checkpoint/stop.sql and save the output to /root/flink/checkpoint/stop.out. After stopping, there must be no dot files left in out1.
  5. Submit /root/flink/checkpoint/run2.sql, which restores from the savepoint (job name flk-seq-2, sink file:///root/flink/checkpoint/out2), put the output in /root/flink/checkpoint/run2.out, and a few seconds later stop through REST while taking a savepoint and save the status response once it has completed to /root/flink/checkpoint/stop2.json.
  6. Submit /root/flink/checkpoint/run3.sql, which runs from the beginning without restoring (job name flk-seq-3, sink file:///root/flink/checkpoint/out3, checkpoint retention RETAIN_ON_CANCELLATION), leave /root/flink/checkpoint/run3.out, and cancel the job after files have been finalized. Save the cancelled job's /jobs/<jid> to /root/flink/checkpoint/job3.json and /jobs/<jid>/checkpoints to /root/flink/checkpoint/checkpoints3.json.
  7. Submit /root/flink/checkpoint/run4.sql (job name flk-seq-4), which restores from the remaining checkpoint and continues writing to the same out3, leave /root/flink/checkpoint/run4.out, and a few seconds later stop through REST while taking a savepoint and save the status response to /root/flink/checkpoint/stop4.json.
  8. In /root/flink/checkpoint/report.json, write savepoint_path, last_id_before_stop, first_id_after_resume, fresh_run_overlap, retained_checkpoint, and leftover_inprogress_files.

Notes

Submit a sequence job with checkpointing on

After flink-up, in /root/flink/checkpoint/run1.sql write a checkpoint interval of 1 s, execution.checkpointing.dir = file:///root/flink/checkpoint/ckpt, execution.checkpointing.savepoint-dir = file:///root/flink/checkpoint/sp, the job name flk-seq-1, a datagen sequence source, and a csv sink and INSERT going to file:///root/flink/checkpoint/out1, and create /root/flink/checkpoint/run1.out with sql-client.sh -f run1.sql > run1.out 2>&1.

Put the SET statements before the INSERT for them to apply to that job. A sequence source is fields.id.kind = sequence, and you give the start and end with fields.id.start and fields.id.end. Keep the rows per second small (20) so that few files are created and it is easy to follow by eye. If you see a Job ID in the output, the job is still running in the background.

Look at the sink directory while it runs

While the job is running, save the output of ls -A /root/flink/checkpoint/out1 to /root/flink/checkpoint/files-running.txt. The finalized part-… and the .part-…inprogress… being written must show up together.

The first part file is finalized only after one checkpoint has completed. With a 1-second interval, waiting a few seconds is enough. Files that start with a dot are invisible without -A, and when measured on this Pod they appear only briefly between checkpoints — take the listing several times at 0.1-second intervals and save it when the two show together. The grader checks whether the finalized files in the listing are still in out1 now.

Get the checkpoint records through REST

Save the /jobs/<jid>/checkpoints of the run1 job to /root/flink/checkpoint/checkpoints.json and /jobs/<jid>/checkpoints/config to /root/flink/checkpoint/checkpoint-config.json. There must be at least one completed checkpoint, and the path of the last one must be this job's ckpt/<jid>/chk-<번호> (the placeholder stands for the number).

The job id is on the Job ID line of run1.out. The latest.completed of the checkpoints response has the number and external_path of the last completed checkpoint, and the config response has the mode (exactly once), the interval (milliseconds), and whether it is retained (externalization).

Take a savepoint and stop

In /root/flink/checkpoint/stop.sql, write the SET for the savepoint directory and STOP JOB '<run1 의 jid>' WITH SAVEPOINT; (the placeholder stands for the jid of run1), and save the output to /root/flink/checkpoint/stop.out. After stopping, there must be no dot files in out1, and the ids of the finalized files must continue from 1 without gaps or duplicates.

STOP JOB is a statement that runs in sql-client and returns the savepoint path as one cell of a table. The path has the shape sp/savepoint--…. The moment the savepoint completes, the file being written is finalized too, so there must be no dot files after stopping — if some remain, it was stopped by cancellation.

Restore from the savepoint and continue writing

Submit /root/flink/checkpoint/run2.sql, with SET 'execution.state-recovery.path' = '<stop.out 의 세이브포인트>'; at the very front (the placeholder stands for the savepoint in stop.out), the job name flk-seq-2, and the sink changed to file:///root/flink/checkpoint/out2, and leave /root/flink/checkpoint/run2.out. A few seconds later, stop while taking a savepoint through REST (POST /jobs/<jid>/stop) and save the status response that has become COMPLETED to /root/flink/checkpoint/stop2.json. The first id of out2 must be right after the last one of out1.

The savepoint also contains, as state, the next number the source will emit. That is why the numbers continue even if you change the sink path. Stopping through REST immediately returns only a request-id, so keep fetching /jobs//savepoints/ until it is COMPLETED and save that.

Run anew without a savepoint and cancel

Submit /root/flink/checkpoint/run3.sql, without a restore path, with the job name flk-seq-3, the sink file:///root/flink/checkpoint/out3, and SET 'execution.checkpointing.externalized-checkpoint-retention' = 'RETAIN_ON_CANCELLATION';, and leave /root/flink/checkpoint/run3.out. After finalized files appear in out3, cancel the job, and save the cancelled job's /jobs/<jid> to /root/flink/checkpoint/job3.json and /jobs/<jid>/checkpoints to /root/flink/checkpoint/checkpoints3.json.

A job that was not restored has its source count again from 1 — the same numbers as out1 come out. Cancellation does not take a savepoint. Without the retention setting, the checkpoint directory is deleted along with the cancellation, so the grader checks whether _metadata remains in the chk directory that checkpoints3.json points to.

Continue writing from the checkpoint you kept

Use the last completed checkpoint in checkpoints3.json as the restore path, submit /root/flink/checkpoint/run4.sql (job name flk-seq-4, with the same sink file:///root/flink/checkpoint/out3), and leave /root/flink/checkpoint/run4.out. A few seconds later, stop while taking a savepoint through REST and save the status response to /root/flink/checkpoint/stop4.json. The finalized files of out3 must continue from 1 without gaps or duplicates.

A checkpoint directory (chk-N) can also be used as a restore path like a savepoint. The new job rewrites from the number at the checkpoint time, and the dot file left half-written at cancellation is not a result. The grader gathers only files that do not start with a dot. Files from two runs with different name tags (uuid) must be mixed together.

Report — recount from the disk

In /root/flink/checkpoint/report.json, write savepoint_path (the savepoint in stop.out), last_id_before_stop (the last id of the finalized files of out1), first_id_after_resume (the first id of out2), fresh_run_overlap (the number of ids finalized into out3 by run3 — the run that started over from 1 — that are also in out1 and out2), retained_checkpoint (the last completed path in checkpoints3.json), and leftover_inprogress_files (the number of dot files left in out3 now).

Every value can be recounted from the disk. A finalized file has a name that starts with part-, and the files of one run share the same name tag (uuid). The files of run3 have the same name tag as the file containing 1. The overlap is the size of the intersection of the two id lists (comm -12).