【问题标题】:How does Flink consume messages from a Kafka topic with multiple partitions, without getting skewed?Flink 如何消费来自具有多个分区的 Kafka 主题的消息,而不会出现偏差?
【发布时间】:2017-08-07 17:36:31
【问题描述】:

假设我们有一个主题的 3 个 kafka 分区,我希望我的事件按小时窗口化,使用事件时间。

当 kafka 消费者在当前窗口之外时,它会停止从分区读取吗?还是打开一个新窗口?如果它正在打开新窗口,那么如果一个分区的事件时间与其他分区相比非常倾斜,那么理论上是否可以让它打开无限数量的窗口并因此耗尽内存?当我们重放一些历史时,这种情况尤其可能发生。

我一直试图从阅读文档中得到这个答案,但找不到关于 Flink 和 Kafka 在分区上的内部结构。一些关于这个特定主题的好的文档将非常受欢迎。

谢谢!

【问题讨论】:

    标签: event-handling apache-kafka apache-flink flink-streaming


    【解决方案1】:

    因此,首先不断读取来自 Kafka 的所有事件,并且进一步的窗口操作对此没有影响。谈到内存不足时,需要考虑更多的事情。

    • 通常您不会存储窗口的每个事件,而只是存储事件的一些聚合
    • 每当关闭窗口时,相应的内存就会被释放。

    更多关于 Kafka 消费者如何与 EventTime 交互的信息(尤其是水印,您可以查看here

    【讨论】:

    • 这非常有用,谢谢。这个例子中的reduce会在窗口累积事件时执行吗?不管Source.windowByEventTime().reduce(someReduceFunc).toSomeSink(foo)
    • 是的,只会存储reduce函数的结果。
    【解决方案2】:

    你可以尝试使用这种风格

    public void runStartFromLatestOffsets() throws Exception {
            // 50 records written to each of 3 partitions before launching a latest-starting consuming job
            final int parallelism = 3;
            final int recordsInEachPartition = 50;
    
            // each partition will be written an extra 200 records
            final int extraRecordsInEachPartition = 200;
    
            // all already existing data in the topic, before the consuming topology has started, should be ignored
            final String topicName = writeSequence("testStartFromLatestOffsetsTopic", recordsInEachPartition, parallelism, 1);
    
            // the committed offsets should be ignored
            KafkaTestEnvironment.KafkaOffsetHandler kafkaOffsetHandler = kafkaServer.createOffsetHandler();
            kafkaOffsetHandler.setCommittedOffset(topicName, 0, 23);
            kafkaOffsetHandler.setCommittedOffset(topicName, 1, 31);
    kafkaOffsetHandler.setCommittedOffset(topicName, 2, 43);
    

    【讨论】:

      猜你喜欢
      • 2020-12-12
      • 2018-12-31
      • 1970-01-01
      • 2023-01-27
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多