【发布时间】:2020-02-20 16:14:01
【问题描述】:
TL;DR:目前在 Flink 中保证事件的事件时间顺序的最佳方案是什么?
我使用 Flink 1.8.0 和 Kafka 2.2.1。我需要通过事件时间戳来保证事件的正确顺序。我每隔 1 秒生成一次周期性水印。我将 FlinkKafkaConsumer 与 AscendingTimestampExtractor 一起使用:
val rawConsumer = new FlinkKafkaConsumer[T](topicName, deserializationSchema, kafkaConsumerConfig)
.assignTimestampsAndWatermarks(new AscendingTimestampExtractor[T] {
override def extractAscendingTimestamp(element: T): Long =
timestampExtractor(element)
})
.addSource(consumer)(deserializationSchema.getProducedType).uid(sourceId).name(sourceId)
然后处理:
myStream
.keyBy(ev => (ev.name, ev.group))
.mapWithState[ResultEvent, ResultEvent](DefaultCalculator.calculateResultEventState)
我意识到,对于在同一毫秒或几毫秒后出现的无序事件,Flink 不会更正顺序。我在文档中找到的内容:
水印触发计算最大时间戳(即 end-timestamp - 1)小于新水印的所有窗口
所以我准备了额外的处理步骤来保证事件时间的顺序:
myStream
.timeWindowAll(Time.milliseconds(100))
.apply((window, input, out: Collector[MyEvent]) => input
.toList.sortBy(_.getTimestamp)
.foreach(out.collect) // this windowing guarantee correct order by event time
)(TypeInformation.of(classOf[MyEvent]))
.keyBy(ev => (ev.name, ev.group))
.mapWithState[ResultEvent, ResultEvent](DefaultScoring.calculateResultEventState)
但是,我觉得这个解决方案很难看,而且看起来像是一种解决方法。我也关注per-partition watermarks of KafkaSource
理想情况下,我想将订单保证放在 KafkaSource 中,并为每个 kafka 分区保留它,就像每个分区的水印一样。有可能这样做吗? 目前Flink中保证事件事件时间顺序的最佳方案是什么?
【问题讨论】:
标签: apache-kafka apache-flink flink-streaming stream-processing