TT Lab
开始
学习 学习路径 课程

Apache Flink — 用真正的引擎跑流处理

启动集群,把一个作业从头跟到尾

在 TT Lab 中继续学习

目标

启动 Flink 本地集群,用 REST API 确认其结构,然后把一个批处理作业从提交一直跟到结束。亲手改动并行度、槽位和失败记录,并把结果整理成报告。

为什么重要

在生产环境中收到的问题——作业为什么停了、并行度调高了为什么没变化、为什么没有重启——关乎的不是 SQL,而是引擎的结构。JobManager 为每个作业创建 JobMaster 并分配槽位,所有记录都通过 REST 公开。本实验的评分器不会查询集群,只读取你保存成文件的 REST 响应和 sql-client 输出,聚合的期望值则直接从源 CSV 计算后再核对。所以响应不要手工修改,要原样保存 curl 的输出。

步骤

  1. 用 flink-up 启动集群,把 /overview 的响应保存到 /root/flink/cluster/overview.json(必须能看到 1 个 TaskManager、2 个槽位)。
  2. 把 /taskmanagers 的响应保存到 /root/flink/cluster/taskmanagers.json,并把其中的 memoryConfiguration 换算成 MiB 并四舍五入,把得到的七个值写入 /root/flink/cluster/memory.json。
  3. 在 /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。
  4. 作业结束后,把 /jobs/overview 的响应保存到 /root/flink/cluster/jobs.json。
  5. 创建用并行度 2、作业名称 flk-first-p2 运行同样聚合的 /root/flink/cluster/p2.sql,把输出保存到 /root/flink/cluster/p2.out,把该作业的 /jobs/<jid> 响应保存到 /root/flink/cluster/p2-job.json。顶点并行度中必须能看到 2。
  6. 把 /opt/flink/conf/config.yaml 中的 taskmanager.numberOfTaskSlots 改为 4 并重新启动集群,然后把 /overview 保存到 /root/flink/cluster/overview-4.json。
  7. 用 /root/flink/cluster/fail.sql 运行一个作业名称为 flk-bad-cast、把 status CAST 成 INT 时会挂掉的作业,并把该作业的 /jobs/<jid> 保存到 /root/flink/cluster/fail-job.json,把 /jobs/<jid>/exceptions 保存到 /root/flink/cluster/fail-exceptions.json。
  8. 在 /root/flink/cluster/report.json 中写入 slots_total、first_job_id、p2_job_id、failed_job_id、restart_strategy、root_cause。

参考

启动集群并获取概览

用 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 中抄过来即可。