每日将Azure Databricks Parquet数据上传至Azure Event Hub如何避免重复?
实现思路
核心通过水印标记机制识别增量数据:记录上一次成功上传到Event Hub的最大数据边界(比如数据生成时间、自增主键ID),每次任务运行时仅读取大于该边界的新数据,上传完成后更新边界值,从逻辑上避免重复上传。
前置准备
- 你的Parquet表需要存在可识别增量的字段:推荐使用
数据生成时间戳、自增唯一ID,如果Parquet已经按dt(日期格式yyyy-MM-dd)分区,可直接用分区字段做增量识别,无需额外水印 - Databricks集群提前安装
azure-eventhub依赖:可在集群库管理界面搜索安装,或在Notebook开头执行%pip install azure-eventhub - 提前准备Azure Event Hub的连接字符串、Event Hub名称
示例代码(基于水印方案)
from azure.eventhub import EventHubProducerClient, EventData from pyspark.sql import functions as F import json # -------------------------- 配置参数 -------------------------- EVENT_HUB_CONN_STR = "你的Event Hub连接字符串" EVENT_HUB_NAME = "你的Event Hub名称" PARQUET_PATH = "你的Parquet文件存储路径" # 水印文件存储路径,用于记录上次上传的最大增量字段值,存在DBFS或Azure Blob均可 WATERMARK_PATH = "/dbfs/mnt/event_hub_upload_watermark.txt" # 你Parquet表里的增量标识字段,这里用数据生成时间戳为例,也可以换成自增ID INCREMENT_COL = "data_create_time" # -------------------------- 1. 读取上次上传的水印 -------------------------- last_watermark = None try: with open(WATERMARK_PATH, "r") as f: last_watermark = f.read().strip() except FileNotFoundError: # 首次运行无水印,默认从最早数据开始同步,可根据需求自定义初始值 last_watermark = "1970-01-01 00:00:00" # -------------------------- 2. 读取本次要上传的增量数据 -------------------------- df = spark.read.parquet(PARQUET_PATH) # 仅筛选大于上次水印的新增数据 incremental_df = df.filter(F.col(INCREMENT_COL) > last_watermark) # 没有新增数据直接退出 if incremental_df.count() == 0: dbutils.notebook.exit("无新增数据,无需上传") # -------------------------- 3. 上传数据到Event Hub -------------------------- producer = EventHubProducerClient.from_connection_string( conn_str=EVENT_HUB_CONN_STR, eventhub_name=EVENT_HUB_NAME ) # 拿本次同步的最大边界值,用于更新水印 current_max_value = incremental_df.agg(F.max(INCREMENT_COL).alias("max_val")).collect()[0]["max_val"] # 转成JSON格式发送,可根据需求调整序列化逻辑 data_list = [json.dumps(row.asDict()) for row in incremental_df.collect()] try: event_data_batch = producer.create_batch() for data in data_list: try: event_data_batch.add(EventData(data)) except ValueError: # 批次满了先发送,再开新批次 producer.send_batch(event_data_batch) event_data_batch = producer.create_batch() event_data_batch.add(EventData(data)) # 发送最后一个批次 producer.send_batch(event_data_batch) finally: producer.close() # -------------------------- 4. 上传成功后更新水印 -------------------------- with open(WATERMARK_PATH, "w") as f: f.write(str(current_max_value))
简化方案(Parquet按日期分区场景)
如果你的Parquet表已经按日期做了分区,且每日调度任务同步当天/前一天的数据,可以不用维护水印,直接读取对应日期的分区即可:
# 以同步前一天数据为例 yesterday = F.date_sub(F.current_date(), 1).cast("string") incremental_df = spark.read.parquet(PARQUET_PATH).filter(F.col("dt") == yesterday) # 后续发送逻辑和上面一致
注意事项
- 水印存储可根据需求替换为Azure SQL、Azure Blob等更高可用的存储介质,避免DBFS文件丢失导致重复同步
- 新增数据量较大时,可调整批量发送的大小,或采用分布式发送方式提升性能
- 代码可增加异常捕获逻辑,确保只有全部数据发送成功后才更新水印,避免数据丢失
- Event Hub默认是至少一次送达语义,如果需要严格不重复,可在消费端根据数据唯一ID做幂等处理
内容的提问来源于stack exchange,提问作者Saswat Ray
相关产品推荐
相关产品推荐

