作者: shenlei
来源: 有赞技术团队
原文发表于: 2021 年 1 月 28 日
原文: https://tech.youzan.com/flink_res_optimize/
背景:把资源配置从经验带回到指标
随着 Flink 集群迁移到 Kubernetes,有赞获得了大促期间弹性扩缩容的便利,也减少了一部分集群运维成本。不过,任务资源仍由用户在实时平台配置;经验不足时,容易出现配置远高于实际需要的情况。例如,某个任务用 4 个并发已经能够满足业务处理需求,却配置了 16 个并发,造成集群资源浪费。
文章把资源分析聚焦在堆内存和 CPU 两方面:一方面从 TaskManager 的 GC 日志估计堆内存需求,另一方面把 Kafka 输入速率与任务中最慢处理环节的能力对照,判断并发度是否可能过高或不足。平台先发现可优化对象,再由管理员与业务方确认,最后调整配置。
Flink 任务所涉及的资源并不只有这两项。作者列出的五类资源是:内存、本地磁盘或云盘、HDFS/S3 等外部存储、HBase/MySQL/Redis 等数据服务、CPU 和网卡。本文的方案主要分析内存与 CPU,不覆盖其他资源是否成为瓶颈。
用 GC 日志估算堆内存
GC 日志记录每次垃圾回收前后堆不同区域的变化。通过日志分析,可以看到 TaskManager 的堆大小、年轻代和老年代分配、GC 次数与停顿时间,以及 Full GC 后老年代剩余空间。文章把 Full GC 后老年代的存活占用记为 M,并依据文中引用的 Java 堆经验规则给出初始估算:堆总量约为 3~4 × M,新生代约为 1~1.5 × M,老年代约为 2~3 × M。对可能出现流量暴涨的任务,实际配置还可以留出更多余量。
如果由此估算出的推荐堆大小与当前分配差距很大,平台可以把任务列为待分析对象。这个估算依赖真实的 GC 行为;它不是适用于所有任务的固定容量公式。
GC 日志收集与分析流程
有赞案例用 GC Viewer 分析 TaskManager GC 日志。文章给出的命令是:
java -jar gcviewer-1.37-SNAPSHOT.jar gc.log summary.csv
其中 gc.log 是一个 TaskManager 的 GC 日志,summary.csv 是分析结果。GC Viewer 需要先获取项目代码并构建后使用;文章没有给出适用于所有读者的部署包或通用日志下载地址。
文中列出的部分分析字段包括:
| 字段 | 含义 |
|---|---|
RunHours |
Flink 任务运行小时数 |
YGSize / YGUsePC |
TaskManager 新生代最大分配量(MB)/ 新生代最大使用率 |
OGSize / OGUsePC |
老年代最大分配量(MB)/ 老年代最大使用率 |
YGCount |
Young GC 次数(原文表格中拼写为 YGCoun) |
YGPerTime |
Young GC 平均停顿时间,单位秒 |
FGCount / FGAllTime |
Full GC 次数 / Full GC 总耗时,单位秒 |
Throughput |
TaskManager 吞吐量(原文拼写为 Throught) |
AVG PT |
平均每次 Young GC 晋升到老年代的对象大小,对应分析结果中的 avgPromotion |
Rec Heap / RecNewHeap / RecOldHeap |
按上述规则估算的推荐总堆、新生代和老年代大小 |
原文建议按 TaskManager 的 Young GC 次数排序,分析次数最多的前 16 个实例;这是一种抽样策略,可能遗漏其他实例上的瓶颈。Yarn 环境中可通过 TaskManager 日志链接查看并经 HTTP 下载。文中的 Kubernetes 环境则先把日志写入 Pod 挂载的云盘(通过 hostPath volume 挂载),再由 Filebeat 监听日志变更并发送到 Kafka;内部日志服务消费这些记录并提供下载接口。实际部署必须使用所在平台允许的、经过权限控制的日志获取方式。
文章案例的环境具有明显的时间和平台范围:内部 JDK 1.8、每个 Kubernetes Pod 限制约 0.6~1 个 CPU core,GC Viewer 包名为 1.37-SNAPSHOT。案例当时选择 Parallel Scavenge 作为年轻代收集器,并在 Serial Old 与 Parallel Old 中选择 Serial Old;作者的考虑是单 core 配额下,多线程老年代回收可能带来线程切换开销。GC 选择应结合吞吐量与延迟权衡:降低单次停顿时间并不必然提高总吞吐量,因为回收次数也可能上升;而过高延迟还可能影响 JobManager、TaskManager、ResourceManager 的心跳。以上是 2021 年内部环境的选择,不应直接当作今天所有集群的建议。
对照 Kafka 输入与处理能力
第二个分析视角是任务能否以合适的资源处理输入数据。文章主要针对 Kafka 数据源:Kafka Topic 的单位时间输入速率可以通过 Kafka Broker 的 JMX 指标获取,也可以从 Flink REST Monitoring API 汇总 Kafka Source Task 的输入;作者选择直接读 Broker 指标,因为反压可能影响 Source 端观察到的输入值。
要估算任务处理能力,需要找到最慢的 Operator 或 Task。一个 Flink 任务的端到端能力会被最慢的一段限制。原文举例:Kafka 输入为 20,000 Record/s,某 Map 算子并发度为 10,每条记录在 Dubbo 调用中的请求往返约为 10 ms,则这个算子的理论处理能力约为:
1000 ms/s ÷ 10 ms/record × 10 并发 = 1000 Record/s
在这个例子里,处理环节的能力显著低于输入速率,算子会成为瓶颈。为定位最慢环节,有赞在 Flink 源码层加入了单条记录处理时间自定义指标 taskOneRecordDealTime,并通过 Flink REST API 读取。平台遍历任务中的 Vertex(可视为 JobGraph 中的 JobVertex),对每个 Vertex 的各个 Task 读取该指标并记录最大值,最后比较各 Vertex,确定处理时间最大的逻辑段。若 Source 到 Sink 全部 chain 在一起,则文章建议着眼于最慢 Operator 的逻辑。
REST 查询的路径格式在原文中写作:
base_flink_web_ui_url/jobs/:jobid
该接口用于取得任务的 Vertex 列表。再对 Vertex 下的 metrics 接口加 ?get= 查询具体指标。例如,查询名为 Filter.numRecordsOut 的指标时,原文示意为:
metrics?get=0.Filter.numRecordsOut
其中 0 表示该 Vertex Task 的 ID。读取有赞添加的单条记录处理时间时,使用 0.taskOneRecordDealTime。多个指标可以在 get 后用逗号分隔。REST 路径和指标名需结合实际 Flink 版本、部署方式以及是否包含该自定义指标确认;这个自定义指标不是所有 Flink 部署自带的通用指标。
用输入、输出与单条处理时间判断并发度
设 Kafka Topic 的单位时间输入为 S,最慢 Task 所在 JobVertex 的并发度为 P,该 JobVertex 单位时间总输出为 O,最慢 Task 的最大单条记录处理时间为 T。文中的估算 1 second / T × P 表示由单条处理耗时和并发度推测的处理能力。作者给出的判断逻辑如下:
- 若
O约等于S,且1 second / T × P远大于S,则考虑降低并发度。 - 若
O约等于S,且1 second / T × P约等于S,则不因这一组指标调整并发度。 - 若
O远小于S,且1 second / T × P也远小于S,则考虑增加并发度。
平台会周期性检测;连续多次符合第一种情况时,向平台管理员告警,提示任务可能为获得当前处理能力配置了过多 CPU。输出与输入吻合时说明当前处理结果没有明显缺口,但降低并发度仍要核对业务目标和其他瓶颈。过滤、聚合、扇出、异步操作、数据倾斜等逻辑都会影响输入、输出与吞吐估算之间的关系,因此这些比较应作为排查线索,而不能单独充当普遍容量公式。
平台实践:自动发现,人工确认
有赞实时平台每天定时扫描正在运行的 Flink 任务。内存方面,它结合 GC 日志与堆大小估算规则,计算推荐堆并与当前分配比较;差距达到平台设定的告警条件时,通知管理员。收到告警后,管理员还会查看消息处理能力。如果最慢 Vertex 的总输出大致跟上 Kafka 输入,而基于并发度和单条处理指标估计出的能力远高于输入,平台会把任务视为可以讨论降低并发度的候选。
告警不是自动改配命令。具体调整幅度由管理员与业务方沟通确定。文章也说明当时流程尚未完全自动化,后半段仍需要人工判断。
结语与适用边界
这套探索的第一步,是用 GC 与处理指标自动筛选可能浪费资源的实时任务,再由平台人员分析并和业务方共同决定是否调整。作者提出的后续方向包括结合任务在不同时段的历史资源使用情况进行自动推测和调配、借鉴离线任务资源优化经验,并探索任务自行弹性扩缩容。
这是针对有赞当时的运行平台、自定义指标和硬件限制形成的实践记录。文章里的 3~4 × M 堆估算、前 16 个 TaskManager 取样、GC 选择以及输入/输出比较都依赖当时的日志质量、任务结构和平台实现。Full GC 后老年代占用与可用空闲空间不是同一个概念;如果任务没有 Full GC,M 不可直接套用。堆分析也不涵盖 Flink 托管或堆外内存、容器总内存限制、磁盘及外部服务。指标需要访问控制;在生产任务调整前,应评估流量峰值、状态恢复和回退方案。文中没有提供可普遍套用的“安全最小堆”或资源节省保证。
来源说明: 本文原文网页未明示文章转载许可;原文信息及代码示例来源于有赞技术团队发布的上述文章。原文案例数据不代表本文作者对其他平台的实测结果。











暂无评论内容