Azure Databricks读取Event Hubs时重复读取数据的求助
问题原因
你当前的代码硬编码了startingEventPosition参数,其中offset: -1表示从Event Hub的最早可用位置开始消费。每次启动流作业时,这个配置会覆盖Spark checkpoint中保存的消费进度,导致作业重启后重新读取所有历史数据,而非从上次停止的位置继续。
解决方案
要实现仅读取新数据、避免重复,核心是让Spark依赖checkpoint跟踪消费位置,而非每次强制指定起始位置:
移除硬编码的起始位置配置
删除代码中关于startingEventPosition的所有内容,包括import json相关代码和ehConf["eventhubs.startingPosition"]这一行。首次运行可选:从最新位置开始消费
如果第一次运行时不想读取历史数据,只想从当前时刻之后的新数据开始,可以在首次运行时临时添加以下配置,运行一次后再移除:import json startingEventPosition = { "offset": "@latest", "seqNo": -1, "enqueuedTime": None, "isInclusive": True } ehConf["eventhubs.startingPosition"] = json.dumps(startingEventPosition)确保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
相关产品推荐
相关产品推荐

