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

每日将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 07:36:05