【发布时间】: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