【发布时间】:2021-09-18 15:35:48
【问题描述】:
我有一个配置了 Kafka 连接器的 Flink 管道。
我已使用以下方法将水印生成频率设置为 2 秒:
env.getConfig().setAutoWatermarkInterval(2000);
现在我的滚动窗口是 60 秒的流窗口,我们在其中进行一些聚合,并且我们根据我们的数据字段之一的时间戳进行基于事件时间的处理。
我没有为我的水印策略或我的信息流配置 allowedLateness。
final ConnectorConfig topicConfig = config.forTopic("mytopic");
final FlinkKafkaConsumer<MyPojo> myEvents = new FlinkKafkaConsumer<>(
topicConfig.name(),
AvroDeserializationSchema.forSpecific(MyPojo.class),
topicConfig.forConsumer()
);
myEvents.setStartFromLatest();
myEvents.assignTimestampsAndWatermarks(
WatermarkStrategy
.<MyPojo>forBoundedOutOfOrderness(
Duration.ofSeconds(30))
.withIdleness(Duration.ofSeconds(120))
.withTimestampAssigner((evt, timestamp) -> evt.event_timestamp_field));
Q.1 从我正在阅读的内容来看,我的时间 0-60 的窗口将在 90 秒后计算,30-90 将在 120 秒后计算,依此类推。然而,由于我们正在做翻转窗口,即没有重叠,我的猜测是没有 30-90 窗口,0-60 之后的下一个窗口是 60-120,在 150 秒标记处触发,对吗?
Q.2 如果没有 allowedLateness,所有迟到的数据都将被丢弃,例如。 90 秒后到达的时间戳为 45 的事件被认为是乱序的,将超出第一个窗口,即 0-60。对于窗口 60-120,事件时间戳不匹配,因此将被丢弃且不包含在窗口在 150 秒处触发,对吗?
问题 3。对于源空闲持续时间,我选择 120 表示如果该主题的任何 Kakfa 分区没有数据,则在 2 分钟后将其标记为空闲,然后发送其他活动分区的水印。我的问题是关于这个数字的选择,即 2 分钟,以及它是否与窗口持续时间(60 秒)或无序(30 秒)有关。如果是这样,我应该在这里记住什么来进行恰当的选择,这样我就不会因为空闲分区导致的非高级水印而导致数据延迟?
或者 120 等待的时间太长,我可能会丢失数据,因此我应该将其设置为远小于 OutOfOrderness 持续时间以确保 0 数据丢失?
编辑:添加更多代码
【问题讨论】:
标签: java apache-kafka apache-flink flink-streaming