Flink FORWARD连接算子数据序列化内部机制及作业图相关问题咨询
Flink作业图相关问题解答
问题1:Transaction Source包含的所有算子是否均通过引用/防御性拷贝的方式共享数据?
是的。如果Transaction Source包含的算子满足算子链合并条件:并行度一致、属于同一个slot共享组、连接方式为FORWARD转发,就会被合并为同一个算子链。链内算子传递数据时默认直接传递对象引用,仅在关闭对象复用配置的场景下才会产生防御性拷贝,全程不会出现序列化或者跨网络传输的开销。
问题2:Transaction Source与Dynamic Partitioning之间的交互逻辑是什么?二者是否因属于FORWARD数据连接而被链化,仅因广播存在才在视觉上呈现分离状态?
二者的连接方式确实是FORWARD数据转发,本身符合算子链的合并要求,数据传递逻辑和同链内算子完全一致,走本地内存传递,要么用对象引用要么仅做防御性拷贝,没有额外序列化、网络开销。
视觉上呈现分离状态的核心原因就是Dynamic Partitioning存在多输入:除了从Transaction Source转发过来的主流,还有一条广播输入流,Flink的算子链规则不允许有多个输入的算子和上游单输入算子合并展示,所以才会在作业图上拆分为两个独立节点,本质上二者的数据交互和已链化的算子没有区别。
内容的提问来源于stack exchange,提问作者Raúl García
相关产品推荐
相关产品推荐

