SevenTnewSAI 与科技新闻,深度解读

流处理 · FLIP-21

Apache Flink 想停止在每个算子处复制数据,而这只是最容易的部分

FLIP-21 将让 Apache Flink 不再在每个链式算子之间复制数据,以三种可选模式取代长期存在的安全默认设置。DEFAULT、COPY_PER_OPERATOR 和 FULL_REUSE 在速度与谨慎之间各取平衡。迁移问题仍未解决。

Emmanuel Fabrice Omgbwa Yasse AI 辅助

2026-09-14 · 阅读需 4 分钟

Apache Flink 想停止在每个算子处复制数据,而这只是最容易的部分

Apache Flink 的流处理运行时中,每当数据在两个链式算子之间传递,运行时都会将其复制。这一行为是有意为之,目的是在状态后端把对象存储在堆上时,避免可变对象带来的麻烦。复制保证了算子向下游传递的内容,不会被其他仍指向同一对象的代码在背后修改。

正在讨论中的变更提案 FLIP-21 的目标是在不放弃这一保证的前提下消除大部分开销。其文档列出了当前做法的四个问题。对 Avro、Thrift 和 JSON 等复杂类型而言,复制的代价高昂。shuffle 之后运行的键控操作从来不需要复制,因为它们位于链的最前端。DataSet API 不在每一步复制,因此 Flink 的两部分对同一份数据遵循不同的规则。而控制整体行为的选项 enableObjectReuse() 名称有误,因为它实际上并不复用对象。

该文档以定性方式描述这些开销,并未给出任何基准测试数字,因此提案没有说明一条典型管道在复制上究竟损失了多少。四项抱怨中有两项关乎开发者而非吞吐量,它们的重要性有一个更不显眼的原因。同一个项目中两套 API 对对象标识的处理方式不同,使人难以推断在作业的某一时点上引用意味着什么。一个名字承诺复用、实际却另有所指的标志更糟,因为它会在“引用能保持有效多久”这一点上植入错误的心智模型。

三种模式,以及控制它们的开关

FLIP-21 没有一刀切地关闭复制,而是定义了三种模式,让用户自行选择在安全与性能之间所处的位置。

模式作用代价
DEFAULT仅在反序列化时创建新对象,之后不再复制直接传递提案将其设为 DataStream 和 DataSet API 的新标准
COPY_PER_OPERATOR在每个算子之间复制数据Flink 当前的行为;安全,但更慢
FULL_REUSE尽可能复用对象最快的选项;对用户代码的谨慎程度要求最高

提案期望 DEFAULT 成为 DataStream 和 DataSet 两套 API 的新标准。它把唯一无法避免的复制限制在反序列化这一个点上,此后对象便可原封不动地传递。另外两种模式分列其两侧。COPY_PER_OPERATOR 是 Flink 今天的做法。FULL_REUSE 把复用推到极致,以速度换取将处理共享引用的负担转移给接收对象的函数。风险并不对等:COPY_PER_OPERATOR 慢但宽容,而 FULL_REUSE 的失败方式取决于周边代码写得有多小心。

与之相关的配置开关 pipeline.object-reuse 默认值为 false。将其设为 true 后,Flink 会在内部复用对象,用于反序列化以及向用户函数传入数据。提案自身关于该开关的说明就是这一权衡的简版:它能减少对象创建和垃圾回收开销,但函数必须正确处理被复用的对象,函数调用返回后引用不保证仍然有效,并且该设置在接近生产环境之前需要经过充分测试。

迁移:彻底切换还是向后兼容

更难的问题是现有管道如何过渡到新的默认设置。目前有两种迁移方案摆在桌面上。社区倾向于第一种,前期改动更大,但之后的结果更干净。

第二个顾虑是那些已经依赖当前复制行为的应用会怎样。提案认为这一群体规模很小,并指出这些用户可以通过配置保留旧模式,因此不将其视为重大障碍。邮件列表上的讨论总体正面。目前尚未设定可用目标或发布日期。

对于今天正在运行 Flink 的团队来说,眼下的变化几乎为零。在新的默认设置到来之前,每个算子都复制的规则依然有效,pipeline.object-reuse 仍默认关闭。安全默认设置所依据的条件早已变化,它的代价却可能持续很久。只有当有人把账算清楚并力主改变时,它才会松动。

每天早晨用 3 分钟掌握科技要闻

每个工作日一封邮件,只讲真正重要的 AI 与科技动态。