启动集群,把一个作业从头跟到尾
目标
启动 Flink 本地集群,用 REST API 确认其结构,然后把一个批处理作业从提交一直跟到结束。亲手改动并行度、槽位和失败记录,并把结果整理成报告。
为什么重要
在生产环境中收到的问题——作业为什么停了、并行度调高了为什么没变化、为什么没有重启——关乎的不是 SQL,而是引擎的结构。JobManager 为每个作业创建 JobMaster 并分配槽位,所有记录都通过 REST 公开。本实验的评分器不会查询集群,只读取你保存成文件的 REST 响应和 sql-client 输出,聚合的期望值则直接从源 CSV 计算后再核对。所以响应不要手工修改,要原样保存 curl 的输出。
步骤
- 用
flink-up启动集群,把/overview的响应保存到 /root/flink/cluster/overview.json(必须能看到 1 个 TaskManager、2 个槽位)。 - 把
/taskmanagers的响应保存到 /root/flink/cluster/taskmanagers.json,并把其中的memoryConfiguration换算成 MiB 并四舍五入,把得到的七个值写入 /root/flink/cluster/memory.json。 - 在 /root/flink/cluster/first.sql 中写一条 SQL:使用批处理模式和作业名称
flk-first-orders,读取/opt/lab/fixtures/data/cluster_orders.csv,按 status 输出orders(订单数)和revenue(amount 之和);并把sql-client.sh -f的输出保存到 /root/flink/cluster/first.out。 - 作业结束后,把
/jobs/overview的响应保存到 /root/flink/cluster/jobs.json。 - 创建用并行度 2、作业名称
flk-first-p2运行同样聚合的 /root/flink/cluster/p2.sql,把输出保存到 /root/flink/cluster/p2.out,把该作业的/jobs/<jid>响应保存到 /root/flink/cluster/p2-job.json。顶点并行度中必须能看到 2。 - 把
/opt/flink/conf/config.yaml中的taskmanager.numberOfTaskSlots改为 4 并重新启动集群,然后把/overview保存到 /root/flink/cluster/overview-4.json。 - 用 /root/flink/cluster/fail.sql 运行一个作业名称为
flk-bad-cast、把statusCAST 成INT时会挂掉的作业,并把该作业的/jobs/<jid>保存到 /root/flink/cluster/fail-job.json,把/jobs/<jid>/exceptions保存到 /root/flink/cluster/fail-exceptions.json。 - 在 /root/flink/cluster/report.json 中写入
slots_total、first_job_id、p2_job_id、failed_job_id、restart_strategy、root_cause。
参考
- 源列:
order_id BIGINT, user_id STRING, status STRING, amount INT, order_time TIMESTAMP(3)(没有表头的 CSV)。 - 运行 SQL 文件:
sql-client.sh -f 파일.sql > 파일.out 2>&1(占位符均为文件名)。遇到出错的语句就会停止,输出文件中会留下[ERROR]。 - 结果被设置为以表格(tableau)模式打印。以批处理方式运行时没有
op列,以流处理方式运行时,最前面会多出op列(+I · -U · +U)。 - 查找作业 ID:
curl -s localhost:8081/jobs/overview | jq -r '.jobs[] | select(.name=="잡이름") | .jid'(占位符为作业名称) - 在配置文件中,槽位数是
taskmanager:下缩进的numberOfTaskSlots:行。只改文件是不会生效的(先flink-down再flink-up)。 - 常见错误:批处理作业的自适应调度器会按数据大小重新确定并行度。文件很小就是 1。
- 这个 Pod 没有互联网。内存上限是 2Gi,所以集群和 SQL 客户端加起来要用 1.5GB 左右。
- 官方文档:Flink Architecture · REST API · TaskManager 内存 · Adaptive Batch · Task Failure Recovery
启动集群并获取概览
用 flink-up 启动集群后,把 curl -s localhost:8081/overview 的响应原样保存到 /root/flink/cluster/overview.json。
flink-up 会调用 start-cluster.sh,并一直等到 TaskManager 连上 JobManager。在连上之前获取的响应里,taskmanagers 显示为 0。请先创建用于保存的目录。
把 TaskManager 内存预算读成数字
把 /taskmanagers 的响应保存到 /root/flink/cluster/taskmanagers.json,并把 memoryConfiguration 中的值换算成 MiB 并四舍五入,把 total_process_mb、total_flink_mb、jvm_metaspace_mb、jvm_overhead_mb、task_heap_mb、managed_mb、network_mb 以整数写入 /root/flink/cluster/memory.json。
值的单位是字节。除以 1048576 后四舍五入。进程总量必须等于 Flink 内存 + 元空间 + 开销——这个恒等式如果不成立,说明哪里抄错了。用 jq 的 round 可以一次生成。
提交一个批处理作业
在 /root/flink/cluster/first.sql 中写入 SET 'execution.runtime-mode' = 'batch';、SET 'pipeline.name' = 'flk-first-orders';、读取源 CSV 的 CREATE TABLE,以及按 status 输出 orders(订单数)和 revenue(amount 之和)的 SELECT,并用 sql-client.sh -f first.sql > first.out 2>&1 生成 /root/flink/cluster/first.out。
filesystem 连接器的 path 是 file:///opt/lab/fixtures/data/cluster_orders.csv,format 是 csv。列名要照抄实验说明中的源列,结果列用 AS orders、AS revenue 起别名。如果输出中出现 op 列,说明是以流处理而不是批处理运行的。
在列表中找到已结束的作业
把 curl -s localhost:8081/jobs/overview 的响应保存到 /root/flink/cluster/jobs.json。名称为 flk-first-orders 的作业必须显示为 FINISHED。
JobManager 会在一段时间内记住已结束的作业(关闭集群后就消失了)。如果看不到名称,请检查 first.sql 中的 pipeline.name 是否位于 SELECT 之前——SET 只对其后提交的作业生效。
用并行度 2 重新运行
创建用 SET 'parallelism.default' = '2'; 和作业名称 flk-first-p2 运行同样聚合的 /root/flink/cluster/p2.sql,把输出保存到 /root/flink/cluster/p2.out,并把该作业的 /jobs/<jid> 响应保存到 /root/flink/cluster/p2-job.json。顶点(vertices)并行度的最大值必须是 2,聚合结果必须与并行度为 1 时相同。
如果并行度设成了 2,顶点却全都是 1,那是批处理作业的自适应调度器根据数据大小(小文件)重新确定了并行度。关闭这种自动决定的配置在 execution.batch.adaptive.auto-parallelism 之下。负责汇总结果的 Sink 保持为 1 是正常的。
把槽位改成 4 个并重新启动
把 /opt/flink/conf/config.yaml 中的 taskmanager.numberOfTaskSlots 改为 4,关闭集群后再重新启动,然后把 /overview 的响应保存到 /root/flink/cluster/overview-4.json。
槽位数是 TaskManager 进程启动时读取一次的值。只改文件然后重新获取概览,得到的仍然是 2。请用 flink-down 关闭,再用 flink-up 启动。也要检查是否有同名的行位于其他块中。
故意让它失败并获取记录
把作业名称设为 flk-bad-cast,用 /root/flink/cluster/fail.sql 运行把 status CAST 成 INT 的 SELECT。把该作业的 /jobs/<jid> 保存到 /root/flink/cluster/fail-job.json,把 /jobs/<jid>/exceptions 保存到 /root/flink/cluster/fail-exceptions.json。
把字符串 'paid' 转换成整数的那一刻,任务会抛出异常,而关闭了检查点的作业不会重启,直接变成 FAILED。sql-client 的输出中也会打出错误,但评分器看的是 JobManager 留下的异常记录(exceptionHistory)。
报告——区分原因与重启策略
在 /root/flink/cluster/report.json 中写入 slots_total(当前集群的槽位数)、first_job_id、p2_job_id、failed_job_id(各作业的 jid)、restart_strategy(最上面的异常语句中阻止重启的策略名称)、root_cause(堆栈中最后一个 Caused by: 的异常类,含包名)。
最上面的异常是以“Recovery is suppressed by …”开头的 JobException。它不是原因,而是“没有重启”这一结果。真正的原因顺着堆栈往下,在最后一个 Caused by 中。jid 从前面保存的 JSON 中抄过来即可。