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

Azure Databricks读取Event Hubs时重复读取数据的求助

问题原因

你当前的代码硬编码了startingEventPosition参数,其中offset: -1表示从Event Hub的最早可用位置开始消费。每次启动流作业时,这个配置会覆盖Spark checkpoint中保存的消费进度,导致作业重启后重新读取所有历史数据,而非从上次停止的位置继续。

解决方案

要实现仅读取新数据、避免重复,核心是让Spark依赖checkpoint跟踪消费位置,而非每次强制指定起始位置:

  1. 移除硬编码的起始位置配置
    删除代码中关于startingEventPosition的所有内容,包括import json相关代码和ehConf["eventhubs.startingPosition"]这一行。

  2. 首次运行可选:从最新位置开始消费
    如果第一次运行时不想读取历史数据,只想从当前时刻之后的新数据开始,可以在首次运行时临时添加以下配置,运行一次后再移除:

    import json
    startingEventPosition = {
      "offset": "@latest",  
      "seqNo": -1,            
      "enqueuedTime": None,   
      "isInclusive": True
    }
    ehConf["eventhubs.startingPosition"] = json.dumps(startingEventPosition)
    
  3. 确保checkpoint目录稳定
    当前代码中的checkpoint位置与存储目录绑定,只要该目录不被删除,Spark就会持续保存消费进度,下次启动时自动从上次结束的位置读取新数据。

修改后的完整代码
connectionString = "Endpoint=sb://ingestionlayer1.servicebus.windows.net/;SharedAccessKeyName=userpolicy2;SharedAccessKey=...=;EntityPath=dataingestion1"

ehConf = {}
ehConf['eventhubs.connectionString'] = sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt(connectionString)
ehConf['eventhubs.consumerGroup'] = "$default"

# Read events from the Event Hub
df_new = spark.readStream.format("eventhubs").options(**ehConf).load()
read_stream = df_new.withColumn("body", df_new["body"].cast("string"))

# Writing to Data lake
saveloc = "/mnt/streamingdatastorage2/datastorage2/Bronze/Table1"
read_stream.writeStream.format("delta").option("checkpointLocation", f"{saveloc}/_checkpoint").start(saveloc)
额外注意事项
  • 不要手动删除checkpoint目录,否则会丢失消费进度,导致作业重启后重新读取所有历史数据。
  • 如果同一Event Hub有多个消费端,确保每个消费端使用不同的consumerGroup,避免互相干扰消费位置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 02:45:33