You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

更换Event Hub及消费组后,Source Connector抛“Invalid Offset”错误求助

问题描述

原有环境部署了指向某命名空间下单分区Event Hub的Event Hub Source Connector,近期迁移至同一命名空间下分区数更多的新Event Hub,并创建了新消费组。但其中一个Connect Worker任务持续失败,报错如下:

Caused by: com.microsoft.azure.eventhubs.impl.AmqpException: The supplied offset '34361193416' is invalid. The last offset in the system is '357208'

其他分区的任务运行正常,仅新Event Hub的0号分区因使用旧偏移量导致失败。尝试修改错误容忍度、添加DLQ主题均未解决问题。

解决方案
  • 重置目标分区的消费组偏移:针对新Event Hub的0号分区和对应新消费组,手动将偏移量重置到有效范围。如果不需要回溯历史数据,直接重置到最新偏移量:
    az eventhubs consumer-group eventhub update --namespace-name <你的命名空间名> --eventhub-name <新Event Hub名称> --name <新消费组名> --partition-id 0 --reset-offset-to-latest
    
    若要保留尽可能多的历史数据,重置到最早可用偏移量:
    az eventhubs consumer-group eventhub update --namespace-name <你的命名空间名> --eventhub-name <新Event Hub名称> --name <新消费组名> --partition-id 0 --reset-offset-to-earliest
    
  • 清理Kafka Connect的偏移存储:如果Connector使用Kafka的connect-offsets主题存储偏移,找到对应分区的旧偏移记录并处理。先通过命令查看偏移主题内容:
    kafka-console-consumer.sh --bootstrap-server <你的Kafka集群地址> --topic connect-offsets --from-beginning --property print.key=true
    
    定位到键格式类似{"name":"<你的Connector名称>","connector":"source","task":<任务ID>,"partition":"<命名空间>/<新Event Hub名>/0/<消费组名>"}的记录,要么用kafka-console-producer将其更新为有效偏移量,要么直接删除该记录(重启Connector后会从指定位置开始消费)。
  • 重启单个失败任务:完成偏移重置后,仅重启对应的Connect Worker任务,避免影响其他正常运行的分区任务,确保数据流转不受大范围影响。

内容的提问来源于stack exchange,提问作者Akhil Kondapaneni

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 00:17:20