【发布时间】:2021-04-27 01:51:09
【问题描述】:
我有一个使用 kafka 的 spark 结构化流应用程序,对于这个应用程序,我想监控消费者滞后。我正在使用以下命令来检查消费者滞后。但是我没有得到 CURRENT-OFFSET ,因此 LAG 也是空白的。这是预期的吗?它适用于其他基于 python 的消费者。
命令
kafka-consumer-groups --bootstrap-server <bootstrap-server>:<port> --describe --all-groups
输出
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
spark-kafka-source-b5e8d872-f727-4ed0-a82c-a3d279647942-407459747-driver-0 my_topic 21 - 5546 - consumer-3-bc651181-fc62-4b1a-abdf-fb3e9d244df8 /<ip-address> consumer-3
spark-kafka-source-b5e8d872-f727-4ed0-a82c-a3d279647942-407459747-driver-0 my_topic 7 - 5129 - consumer-3-bc651181-fc62-4b1a-abdf-fb3e9d244df8 /<ip-address> consumer-3
spark-kafka-source-b5e8d872-f727-4ed0-a82c-a3d279647942-407459747-driver-0 my_topic 3 - 5178 - consumer-3-bc651181-fc62-4b1a-abdf-fb3e9d244df8 /<ip-address> consumer-3
spark-kafka-source-b5e8d872-f727-4ed0-a82c-a3d279647942-407459747-driver-0 my_topic 9 - 4969 - consumer-3-bc651181-fc62-4b1a-abdf-fb3e9d244df8 /<ip-address> consumer-3
spark-kafka-source-b5e8d872-f727-4ed0-a82c-a3d279647942-407459747-driver-0 my_topic 2 - 5443 - consumer-3-bc651181-fc62-4b1a-abdf-fb3e9d244df8 /<ip-address> consumer-3
spark-kafka-source-b5e8d872-f727-4ed0-a82c-a3d279647942-407459747-driver-0 my_topic 15 - 5312 - consumer-3-bc651181-fc62-4b1a-abdf-fb3e9d244df8 /<ip-address> consumer-3
【问题讨论】:
-
我正在维护使用“自定义”消费者组 ID 进行提交的项目 (github.com/HeartSaVioR/spark-sql-kafka-offset-committer)。 (这样您就可以将偏移量提交给您想要的组,而不是唯一生成的组。)正如下面的答案所解释的,Spark 没有,并且您不应该尝试修改行为 Spark 完全处理偏移量信息并且不不要把它交给卡夫卡。
标签: apache-spark apache-kafka kafka-consumer-api spark-structured-streaming spark-kafka-integration