如何扩容优化消费超大Kafka Topic的Flink流作业以降低高延迟
优化方案
1. 替换双流Join方案为广播流关联(收益最高)
当前采用双流keyBy后connect的方案需要对每秒30~50万条的大流量Input1做全量shuffle,是延迟高的核心原因。由于Input2属于低频率更新、全量数据量极小的维度数据,完全符合广播流的适用场景:
- 将Input2数据转换为广播流,配置
MapStateDescriptor存储全量维度数据 - Input1流直接和广播流关联,无需shuffle,所有关联逻辑在Input1的本地并行实例完成,省去海量数据的网络传输开销
- 代码示例:
// 定义广播状态描述符 val broadcastStateDescriptor = new MapStateDescriptor[Long, Input2]("input2BroadcastState", classOf[Long], classOf[Input2]) val broadcastStream2 = stream2.broadcast(broadcastStateDescriptor) // 直接用Input1的消费并行度关联,无需额外shuffle stream1 .connect(broadcastStream2) .process(new BroadcastProcessFunction[Input1, Input2, Output] { override def processElement(value: Input1, ctx: BroadcastProcessFunction[Input1, Input2, Output]#ReadOnlyContext, out: Collector[Output]): Unit = { val state = ctx.getBroadcastState(broadcastStateDescriptor) val dim = state.get(value.getKey1) // 组装输出即可 } override def processBroadcastElement(value: Input2, ctx: BroadcastProcessFunction[Input1, Input2, Output]#Context, out: Collector[Output]): Unit = { val state = ctx.getBroadcastState(broadcastStateDescriptor) state.put(value.getKey2, value) } }) // 并行度和Input1消费并行度一致设为120即可,无需用到700这么高的并行度
2. 现有双流Join方案的优化(如果暂不换广播流)
如果要保留现有CoMapFunction方案,可以做以下优化:
解决数据倾斜
查看Flink UI中关联算子各subtask的处理数据量,若出现部分subtask处理数据量远高于均值的情况,说明存在数据倾斜:- 给Input1的关联key加上随机前缀(比如0~9的随机数),将大key打散到10个并行度处理
- 将Input2的每条数据扩容为10份,每份带上对应前缀,和打散后的Input1关联
- 关联后去掉前缀输出即可
增加本地缓存减少状态访问开销
Input2数据更新频率极低,可以在算子内增加Guava本地缓存,缓存热点key的维度数据,避免每次请求都访问MapState:// open方法中初始化 private var localCache: LoadingCache[Long, Input2] = _ override def open(parameters: Configuration): Unit = { // 原有状态初始化逻辑不变 localCache = CacheBuilder.newBuilder() .maximumSize(10000) // 根据热点key数量调整 .expireAfterWrite(1, TimeUnit.HOURS) // 缓存过期时间可设远小于状态TTL .build( new CacheLoader[Long, Input2]() { override def load(key: Long): Input2 = { input2State.get(key) } } ) } // map1中直接查localCache即可提前过滤冗余字段
Input1单条有600个字段,业务仅需100个,在消费到Input1之后就把无用字段删除,减少序列化、反序列化、网络传输的开销
3. 配置优化
状态后端优化
当前taskManagerMemoryManagedSize配置为1b,相当于禁用了托管内存,如果你用HashMapStateBackend,所有状态都存在TM堆内,很容易触发频繁GC:- 建议换RocksDB状态后端,将状态存在堆外,大幅降低GC压力
- 调整
taskManagerMemoryManagedSize为至少2G,满足RocksDB的内存需求
上下游并行度匹配
输出Kafka Topic分区数为120,当前如果输出并行度为700,会导致多个并发写同一个Kafka分区,降低写入效率,建议将输出算子并行度设置为120和分区数匹配开启Checkpoint
当前未开启Checkpoint,一旦任务故障会丢失所有状态,需要等次日Input2数据重新写入才能恢复补全逻辑,建议开启Checkpoint,间隔设为10~30分钟即可调整水印间隔
当前水印间隔为300000ms(5分钟),如果你未用到窗口逻辑可以调大到10分钟,减少水印传输的开销,如果用到窗口逻辑可以根据业务需求调小
内容的提问来源于stack exchange,提问作者user17080223
相关产品推荐
相关产品推荐

