TT Lab
Get started
Learn Learning paths Courses

Apache Flink — Running Streams on a Real Engine

Start a Cluster and Follow One Job to the End

Continue in TT Lab

Goal

Start a local Flink cluster, check its structure through the REST API, and then follow one batch job from submission to completion. Change the parallelism, the slots, and the failure record yourself, and write up the results in a report.

Why it matters

The questions you get in operations — why did the job stop, why is nothing different after I raised the parallelism, why didn't it restart — are about the structure of the engine, not about SQL. The JobManager creates a JobMaster for each job, hands out slots, and exposes every record through REST. The grader of this lab does not ask the cluster anything. It reads only the REST responses and the sql-client output that you saved as files, and compares the aggregation's expected values with values it computes directly from the source CSV. So do not edit the responses by hand; save the curl output exactly as it is.

Steps

  1. Start the cluster with flink-up and save the /overview response to /root/flink/cluster/overview.json (it must show 1 TaskManager and 2 slots).
  2. Save the /taskmanagers response to /root/flink/cluster/taskmanagers.json, and write seven values, its memoryConfiguration converted to MiB and rounded, to /root/flink/cluster/memory.json.
  3. In /root/flink/cluster/first.sql, write SQL that, in batch mode with the job name flk-first-orders, reads /opt/lab/fixtures/data/cluster_orders.csv and produces orders (count) and revenue (sum of amount) for each status, and save the output of sql-client.sh -f to /root/flink/cluster/first.out.
  4. After the job finishes, save the /jobs/overview response to /root/flink/cluster/jobs.json.
  5. Create /root/flink/cluster/p2.sql, which runs the same aggregation with parallelism 2 and the job name flk-first-p2, save its output to /root/flink/cluster/p2.out, and save that job's /jobs/<jid> response to /root/flink/cluster/p2-job.json. A 2 must appear in the vertex parallelism.
  6. Change taskmanager.numberOfTaskSlots in /opt/flink/conf/config.yaml to 4, start the cluster again, and save /overview to /root/flink/cluster/overview-4.json.
  7. With the job name flk-bad-cast, run a job that dies while casting status to INT as /root/flink/cluster/fail.sql, and save that job's /jobs/<jid> to /root/flink/cluster/fail-job.json and /jobs/<jid>/exceptions to /root/flink/cluster/fail-exceptions.json.
  8. In /root/flink/cluster/report.json, write slots_total, first_job_id, p2_job_id, failed_job_id, restart_strategy, and root_cause.

Notes

Start the cluster and get the overview

Start the cluster with flink-up, then save the curl -s localhost:8081/overview response as it is to /root/flink/cluster/overview.json.

flink-up calls start-cluster.sh and waits until the TaskManager attaches to the JobManager. A response received before it attaches shows taskmanagers as 0. Create the directory to save to first.

Read the TaskManager memory budget as numbers

Save the /taskmanagers response to /root/flink/cluster/taskmanagers.json, and write total_process_mb, total_flink_mb, jvm_metaspace_mb, jvm_overhead_mb, task_heap_mb, managed_mb, and network_mb, the values of memoryConfiguration converted to MiB and rounded, as integers to /root/flink/cluster/memory.json.

The values are in bytes. Divide by 1048576 and round. The whole process must equal Flink memory + metaspace + overhead — if this identity does not hold, you copied something wrong somewhere. If you use jq's round, you can build it in one go.

Submit a batch job

In /root/flink/cluster/first.sql, write SET 'execution.runtime-mode' = 'batch';, SET 'pipeline.name' = 'flk-first-orders';, a CREATE TABLE that reads the source CSV, and a SELECT that produces orders (count) and revenue (sum of amount) for each status, and create /root/flink/cluster/first.out with sql-client.sh -f first.sql > first.out 2>&1.

For the filesystem connector, the path is file:///opt/lab/fixtures/data/cluster_orders.csv and the format is csv. Use the source columns from the instructions exactly as they are for the column names, and give the result columns the aliases AS orders and AS revenue. If you see an op column in the output, it ran as streaming, not batch.

Find the finished job in the list

Save the curl -s localhost:8081/jobs/overview response to /root/flink/cluster/jobs.json. The job named flk-first-orders must show as FINISHED.

The JobManager remembers finished jobs for a while (they disappear when you shut down the cluster). If you do not see the name, check that the pipeline.name in first.sql comes before the SELECT — a SET applies from jobs submitted after it.

Run it again with parallelism 2

Create /root/flink/cluster/p2.sql, which runs the same aggregation with SET 'parallelism.default' = '2'; and the job name flk-first-p2, save its output to /root/flink/cluster/p2.out, and save that job's /jobs/<jid> response to /root/flink/cluster/p2-job.json. The maximum of the vertices parallelism must be 2, and the aggregation result must be the same as with parallelism 1.

If you set the parallelism to 2 and all the vertices are 1, the batch job's adaptive scheduler looked at the data size (a small file) and reset the parallelism. The setting that turns off that automatic decision is under execution.batch.adaptive.auto-parallelism. It is normal for the sink that collects the results to stay at 1.

Change the slots to 4 and start again

Change taskmanager.numberOfTaskSlots in /opt/flink/conf/config.yaml to 4, shut down the cluster and start it again, and save the /overview response to /root/flink/cluster/overview-4.json.

The slot count is a value the TaskManager process reads once when it starts. If you only edit the file and fetch the overview again, you still get 2. Shut down with flink-down and start with flink-up. Also check that a line with the same name is not in another block.

Fail it on purpose and get the record

With the job name flk-bad-cast, run a SELECT that casts status to INT as /root/flink/cluster/fail.sql. Save that job's /jobs/<jid> to /root/flink/cluster/fail-job.json and /jobs/<jid>/exceptions to /root/flink/cluster/fail-exceptions.json.

The moment the string 'paid' is converted to an integer, the task throws an exception, and a job with checkpointing off does not restart and goes straight to FAILED. The error is also printed in the sql-client output, but the grader looks at the exception record (exceptionHistory) that the JobManager kept.

Report — separate the cause from the restart strategy

In /root/flink/cluster/report.json, write slots_total (the number of slots of the cluster now), first_job_id, p2_job_id, and failed_job_id (the jid of each job), restart_strategy (the name of the strategy that blocked the restart, from the top exception sentence), and root_cause (the exception class of the last Caused by: in the stack, including the package).

The top exception is a JobException that starts with 'Recovery is suppressed by …'. That is not the cause but the result "it did not restart". The real cause is in the last Caused by, following the stack down. You can copy the jids from the JSON files you saved earlier.