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

如何通过Auto Loader与Event Grid将Event Hub数据接入Databricks?

使用Event Grid + Auto Loader将Event Hub数据摄入Databricks

前置确认

  • 你的Event Grid订阅已设置为监听Event Hub的Microsoft.EventHub.CaptureFileCreated事件(Auto Loader基于文件处理,需要Event Hub先把数据捕获到存储介质)
  • Databricks集群有权限访问Event Hub捕获数据的存储账户(Azure Blob/ADLS Gen2)

步骤1:确保Event Hub已开启捕获

如果还没配置Event Hub捕获,先完成这一步:

  • 把Event Hub的捕获目标指向你的Azure存储账户,文件格式选Avro或JSON(Avro更推荐,性能更好)
  • 设置捕获的时间/大小阈值,确保数据会定期生成文件到存储路径

步骤2:把Event Grid订阅和Databricks作业绑定

Event Grid需要触发Databricks的Auto Loader作业,按以下操作:

  1. 在Databricks里创建一个作业,作业任务就是后面要写的Auto Loader代码
  2. 生成Databricks个人访问令牌(PAT),权限至少要有jobs.runNow
  3. 回到Event Grid订阅的配置页,目标类型选Webhook,填入调用Databricks作业的URL:
    https://<你的Databricks实例域名>/api/2.0/jobs/run-now?job_id=<你的作业ID>
    
  4. 在Webhook头部添加Authorization: Bearer <你的PAT令牌>
  5. 测试订阅触发,确认Event Grid能成功调用Databricks作业

步骤3:编写Auto Loader核心代码

在Databricks作业任务里,写以下代码实现增量读取:

Python示例

# 替换为你的Event Hub捕获存储路径
input_path = "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/eventhub-capture/"

# 替换为你的Delta表存储路径和表名
output_table_path = "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/delta-tables/eventhub-raw/"
output_table_name = "default.eventhub_raw_data"

# 初始化Auto Loader流
df = (spark.readStream
      .format("cloudFiles")
      .option("cloudFiles.format", "avro")  # 和Event Hub捕获的格式保持一致
      .option("cloudFiles.schemaLocation", f"{output_table_path}/_schema")  # 自动保存/推断Schema
      .load(input_path))

# 写入Delta Lake表(支持增量、断点续传)
query = (df.writeStream
         .format("delta")
         .option("checkpointLocation", f"{output_table_path}/_checkpoint")  # 保存流状态
         .option("mergeSchema", "true")  # 自动合并新增字段
         .table(output_table_name))

# 等待流执行完成(作业场景下可根据需求调整)
query.awaitTermination()

关键参数说明

  • cloudFiles.format: 必须和Event Hub捕获的文件格式匹配
  • schemaLocation: 自动保存Schema,避免每次启动作业重复推断
  • checkpointLocation: 记录流的处理状态,确保中断后能从断点继续

步骤4:验证数据摄入

  • 给Event Hub发送测试数据,等待捕获生成文件
  • 检查Event Grid是否触发了Databricks作业
  • 在Databricks里查询目标表,确认数据已入库:
SELECT * FROM default.eventhub_raw_data LIMIT 10;

常见问题排查

  • Event Grid调用作业失败:检查PAT权限、作业ID是否正确,若Databricks是私有网络,需配置防火墙允许Event Grid访问
  • Auto Loader读不到文件:确认存储路径正确,集群有存储账户的访问权限(SAS/MSI均可)
  • Schema变更丢失:确保开启了mergeSchema参数,支持自动合并新增字段

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 06:55:48