Flink作业多数据流共享同一Kafka源和Sink是否安全?相关疑问
关于你提出的几个问题,解答如下:
1. 该fan-out->fan-in架构的影响
- 故障恢复与Checkpoint影响
这种同源多分支再union的架构本身是Flink原生支持的模式,和常规流处理任务的Checkpoint、故障恢复逻辑没有本质差异:Checkpoint的barrier会从源节点同步发送到所有分支,只要各分支的处理延迟差异不大,union阶段的barrier对齐不会出现异常,故障恢复时各分支的状态也会从Checkpoint中正常加载。
可能的负面影响只有两个场景:一是各分支处理延迟差距过大时,union算子需要等最慢的分支的barrier到达才能完成对齐,会拉高Checkpoint的整体耗时,严重时可能触发Checkpoint超时;二是分支数量较多时,全量任务的状态总大小会高于单流处理方案,Checkpoint的存储开销、故障恢复时的状态加载耗时会对应上升。 - 运行性能影响
源节点的数据会被复制N份(N等于分支数)发送到各分支,所以源节点之后的网络传输、数据序列化开销会线性增长,这是该架构主要的性能损耗点。但如果每个分支的第一逻辑就是高过滤比的过滤操作,过滤后的数据量很小,后续的窗口处理、写sink的开销反而会比揉合所有逻辑的单流方案更低,同时架构的可维护性、分支逻辑的可扩展性也会好很多。如果开启Flink算子链优化,各分支内部的无状态算子(filter、flatmap等)会被链在一起执行,不会产生额外的网络开销。
2. 同类方案的落地情况
这类架构在实时多规则处理场景的落地非常普遍,比如实时风控的多规则并行检测、实时用户标签的多维度并行计算、实时数仓的多指标并行统计等场景都有大量生产落地案例。我接触过的同类型业务最高承载过单任务15万TPS的入流量,3个分支并行处理,稳定运行超过1年没有出现过架构层面的稳定性问题,只要资源配置匹配流量规模即可正常运行。
3. Watermark配置时机
必须在过滤操作之前配置水位线策略,最优的做法是在Kafka源节点之后直接生成水位线,再将带水位线的流拆分到各个分支做后续处理。
如果先过滤再生成水位线,一旦过滤操作丢弃了时间戳最新的事件,会直接导致水位线无法正常推进,出现窗口无法触发、数据无限延迟的问题。而且在源之后统一生成水位线,可以保证所有分支的水位线进度一致,避免分支间水位线不同步导致union后窗口计算异常。
内容的提问来源于stack exchange,提问作者Alexis Jalomo
相关产品推荐
相关产品推荐

