【问题标题】:Spark Structured Streaming - kafka offset handlingSpark Structured Streaming - kafka 偏移处理
【发布时间】:2021-05-22 15:07:24
【问题描述】:

当我从最新的偏移量启动我的 Spark Structured Streaming 3.0.1 应用程序时,它运行良好。但是当我想从一些较早的偏移量开始时——例如:

  • startingOffsets 为“最早”
  • startingOffsets 到特定的偏移量,例如 {"MyTopic-v1":{"0":1686734237}}

我可以在日志中看到起始偏移量被正确拾取,但随后发生了一系列搜索(从我定义的位置开始),直到它达到当前的最新偏移量。

我删除了我的检查点目录并尝试了几个选项,但情况始终相同 - 它报告了正确的起始偏移量,但随后需要很长时间才能慢慢寻找最新的并开始处理 - 知道为什么和什么我还要检查吗?

2021-02-19 14:52:23 INFO  KafkaConsumer:1564 - [...] Seeking to offset 1786734237 for partition MyTopic-v1-0
2021-02-19 14:52:23 INFO  KafkaConsumer:1564 - [...] Seeking to offset 1786734737 for partition MyTopic-v1-0
2021-02-19 14:52:23 INFO  KafkaConsumer:1564 - [...] Seeking to offset 1786735237 for partition MyTopic-v1-0
2021-02-19 14:52:23 INFO  KafkaConsumer:1564 - [...] Seeking to offset 1786735737 for partition MyTopic-v1-0
2021-02-19 14:52:23 INFO  KafkaConsumer:1564 - [...] Seeking to offset 1786736237 for partition MyTopic-v1-0
2021-02-19 14:52:23 INFO  KafkaConsumer:1564 - [...] Seeking to offset 1786736737 for partition MyTopic-v1-0
2021-02-19 14:52:23 INFO  KafkaConsumer:1564 - [...] Seeking to offset 1786737237 for partition MyTopic-v1-0

我让应用程序运行了更长的时间,它最终开始生成文件,但我的 100 秒处理触发器没有得到满足,数据显示要晚得多 - 20-30 分钟后。

(我也在 spark 2.4.5 上测试过——同样的问题——也许是一些 kafka 配置?)

【问题讨论】:

  • 会不会是“需要很长时间才能慢慢找到最新的”,但它实际上是在处理数据?否则,像您一样使用startingOffsets 没有任何问题。也许分享一个minimal reproducible example
  • 这实际上是可能的 - 我让应用程序运行了更长的时间,它最终开始生成文件,但我的 100 秒处理触发器没有得到满足,数据显示得更晚 - 在 20 之后 - 30分钟
  • 谢谢@mike!文档对此参数不太清楚 - 我想这可以更正,因为 maxOffsetsPerTrigger 被描述为更像“速率限制”
  • @mike 请在此处发布您的答案,我会将其标记为正确

标签: apache-spark apache-kafka spark-structured-streaming


【解决方案1】:

将选项 startingOffsets 与您展示的 JSON 对象一起使用应该可以正常工作。

您观察到的是,在应用程序第一次启动时,结构化流式处理作业将读取从提供的 (1686734237) 到主题中最后一个可用偏移量的所有(!)偏移量。由于这可能是相当多的消息,因此处理该大块将使第一个微批处理非常繁忙。

请记住,Trigger 选项只是定义了微批处理的触发频率。您应该确保将此触发率与预期的处理时间保持一致。我在这里看到基本上有两个选择:

  • 使用选项maxOffsetsPerTriger 来限制每个触发器/微批处理从 Kafka 获取的偏移量
  • 避免使用任何触发器,因为默认情况下,这将允许您的流在前一个触发器完成数据处理后立即触发

【讨论】:

    猜你喜欢
    • 2020-09-03
    • 2017-02-06
    • 2017-12-31
    • 2019-07-30
    • 2018-09-22
    • 2018-04-26
    • 2017-06-22
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多