Apache Flink — Running Streams on a Real Engine
Start a Cluster and Follow One Job to the End
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
- Start the cluster with
flink-upand save the/overviewresponse to /root/flink/cluster/overview.json (it must show 1 TaskManager and 2 slots). - Save the
/taskmanagersresponse to /root/flink/cluster/taskmanagers.json, and write seven values, itsmemoryConfigurationconverted to MiB and rounded, to /root/flink/cluster/memory.json. - 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.csvand producesorders(count) andrevenue(sum of amount) for each status, and save the output ofsql-client.sh -fto /root/flink/cluster/first.out. - After the job finishes, save the
/jobs/overviewresponse to /root/flink/cluster/jobs.json. - 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. - Change
taskmanager.numberOfTaskSlotsin/opt/flink/conf/config.yamlto 4, start the cluster again, and save/overviewto /root/flink/cluster/overview-4.json. - With the job name
flk-bad-cast, run a job that dies while castingstatustoINTas /root/flink/cluster/fail.sql, and save that job's/jobs/<jid>to /root/flink/cluster/fail-job.json and/jobs/<jid>/exceptionsto /root/flink/cluster/fail-exceptions.json. - In /root/flink/cluster/report.json, write
slots_total,first_job_id,p2_job_id,failed_job_id,restart_strategy, androot_cause.
Notes
- Source columns:
order_id BIGINT, user_id STRING, status STRING, amount INT, order_time TIMESTAMP(3)(a CSV with no header). - Run an SQL file with
sql-client.sh -f 파일.sql > 파일.out 2>&1(the placeholders stand for the file names). It stops at the statement that errors, and[ERROR]remains in the output file. - Results are set to print in tableau mode. When run as batch there is no
opcolumn, and when run as streaming anopcolumn (+I · -U · +U) is added at the front. - Finding a job id:
curl -s localhost:8081/jobs/overview | jq -r '.jobs[] | select(.name=="잡이름") | .jid'(the placeholder stands for the job name) - In the configuration file, the slot count is the
numberOfTaskSlots:line indented undertaskmanager:. Editing only the file does not take effect (flink-down, thenflink-up). - A common mistake: for batch jobs, the adaptive scheduler resets the parallelism based on data size. With a small file you get 1.
- This Pod has no internet. The memory limit is 2Gi, so the cluster and the SQL client together use around 1.5 GB.
- Official docs: Flink Architecture · REST API · TaskManager memory · Adaptive Batch · Task Failure Recovery
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.