Flink对接Kafka时KeyBy与reinterpretAsKeyedStream使用及性能疑问
Flink 1.13流处理作业相关问题解答
一、reinterpretAsKeyedStream()的记录分发逻辑与报错原因
reinterpretAsKeyedStream()是零网络开销的类型转换方法,本身不会执行任何shuffle、重分区逻辑,仅仅是将现有DataStream标记为KeyedStream,完全沿用上游算子的记录分发结果。
这个方法有严格的使用前提:上游输出的所有记录,必须已经严格匹配Flink KeyedStream的路由规则——也就是每个key对应的记录,必须恰好发送到持有该key对应key-group的下游算子实例上,否则就会抛出key-group不匹配异常。
你遇到并行度大于1时报错的核心原因:
- Kafka侧默认分区器的String字段哈希逻辑,和Flink内部key-group的哈希路由逻辑完全不一致,Kafka分区和Flink key-group之间没有对应关系
- 你的Kafka主题有24个分区,作业并行度为4时Flink默认最大并行度为128,共划分128个key-group,每个并行子任务负责32个连续的key-group。报错中实例负责的key-group范围是96~127,却收到了属于85号key-group的记录,自然触发校验失败
- 并行度设为1时所有key-group都由同一个子任务持有,不存在key-group路由错误的可能,所以可以正常运行
你之前的用法是错误的,reinterpretAsKeyedStream()仅适用于上游本身就是按照Flink keyBy规则完成分区的场景(比如上游是keyedStream经过无状态forward算子输出的流),Kafka分区的流不满足这个前提,必须用keyBy()做重分区。
最初的错误代码逻辑如下:
DataStream<Envelope> messageStream = env.addSource(kafkaSource); DataStreamUtils.reinterpretAsKeyedStream(messageStream, Envelope::getId) .process(new EnvelopeMapper(parameters)) .addSink(kafkaSink);
用到的有状态处理函数定义:
public class EnvelopeMapper extends KeyedProcessFunction<String, Envelope, Envelope> { ... }
抛出的核心异常:
java.lang.IllegalArgumentException: KeyGroupRange{startKeyGroup=96, endKeyGroup=127} does not contain key group 85
二、keyBy()放置位置的性能差异
你修改后的代码把keyBy()放在了两个无状态算子之后,和放在source读取后立刻调用的位置相比,性能差异完全取决于两个无状态算子对数据的处理逻辑,核心影响点是shuffle阶段跨网络传输的数据量:
- 如果把
keyBy()放在source读取后(标注Line x的位置):shuffle会在读取Kafka原始数据后立刻执行,两个无状态算子会在shuffle后的下游节点本地运行,跨网络传输的是Kafka读出的原始数据 - 如果把
keyBy()放在两个无状态算子之后、EnvelopeMapper之前:source和两个无状态算子会组成算子链在Kafka消费节点本地运行,不需要网络传输,跨网络shuffle传输的是经过两个无状态算子处理后的数据
实际性能判断规则:
- 如果两个无状态算子会过滤掉大量无效数据、或者裁剪掉大体积无用字段,最终输出数据比原始Kafka数据小很多:
keyBy()放在无状态算子之后性能更好,能大幅降低网络传输开销 - 如果两个无状态算子会扩充字段、增大单条记录体积:
keyBy()放在source之后立刻调用性能更好 - 如果两个无状态算子几乎不改变单条数据体积、也不过滤数据:两种写法性能差异极小,几乎可以忽略
修改后的正确代码逻辑:
DataStream<Envelope> messageStream = env.addSource(kafkaSource) // Line x位置 .map(statelessMapper1) .flatMap(statelessMapper2); messageStream.keyBy(Envelope::getId) .process(new EnvelopeMapper(parameters)) .addSink(kafkaSink);
内容的提问来源于stack exchange,提问作者Ahmed A
相关产品推荐
相关产品推荐

