Flink多Kafka源Union后水印策略及处理时间窗口调优求助
Flink多Kafka源业务处理问题
我是Flink新手,当前有5个Schema不同的无界Kafka数据源,业务要求先对各源消息去重,再按相同key执行外连接。现有实现是通过Union算子合并所有数据源后,使用ProcessWindowFunction转换为统一的CommonObj对象下发至下游。
遇到的问题
问题1:事件时间窗口下消息被丢弃
采用事件时间窗口时,某一个Kafka源的水印始终滞后于其他源,导致消息被丢弃;即便已经设置了1分钟的空闲超时,该问题仍然存在。
问题2:处理时间窗口下高背压与Checkpoint超时
切换为处理时间窗口(无水印)后,所有Kafka源出现高背压,且Checkpoint在10分钟后超时,调整窗口大小(10-30秒)也无法解决该问题。
相关代码
CommonObj实体类
class CommonObj { var id: Long = 0 var entityType: String? = null var timestamp: Long = System.currentTimeMillis() val eventMetas: MutableList<EventMeta> = mutableListOf() var kafkaStreamValue1: KafkaStreamObj1? = null var kafkaStreamValue2: KafkaStreamObj2? = null var kafkaStreamValue3: KafkaStreamObj3? = null var kafkaStreamValue4: KafkaStreamObj4? = null fun buildSinkObj(): SinkObj = .... }
单Kafka源处理代码示例
val watermarkStrategy = WatermarkStrategy.forMonotonousTimestamps<KafkaStreamObj1>() .withIdleness(Duration.ofMinutes(1)) val sourceStream1 = env.fromSource( getKafkaStream1(params), watermarkStrategy, "Kafka Source 1" ) val kafkaSource1 = sourceStream1 .filter { it != null } .map { EventObj<KafkaStreamObj1>( it.id.toString() + it.entity, // 作为分组key it, // 原始对象 it.sequence, // 事件时间戳 mutableListOf(EventMeta(it.transactionId, it.type, it.sequence, it.changed)) ) } .returns(TypeInformation.of(object : TypeHint<EventObj<KafkaStreamObj1>>() {})) .keyBy { it.key } .window(TumblingEventTimeWindows.of(Time.milliseconds(10000))) .reduce { v1, v2 -> if (v1.obj.sequence > v2.obj.sequence) { v1.eventMetaList.addAll(v2.eventMetaList) v1 } else { v2.eventMetaList.addAll(v1.eventMetaList) v2 } } .map { val commonObj = CommonObj() commonObj.id = it.obj.id commonObj.entityType = it.obj.entity commonObj.timestamp = System.currentTimeMillis() commonObj.eventMetas.addAll(it.eventMetaList) commonObj.kafkaStreamValue1 = it.obj.entity commonObj } .returns(TypeInformation.of(object : TypeHint<CommonObj>() {})) return kafkaSource1
Union合并与下游处理代码
kafkaStream1.union(kafkaStream2,kafkaStream3,kafkaStream4,kafkaStream5) .keyBy(keySelector) .window(TumblingEventTimeWindows.of(Time.milliseconds(10000))) .process(EventProcessFunction(params)) .sinkTo(kafkaSink())
内容的提问来源于stack exchange,提问作者watson
相关产品推荐
相关产品推荐

