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

如何使用Scala在Databricks集群实现Eventhub到ADLS Gen2的流数据摄入

Azure Event Hub 流数据摄入ADLS Gen2 实现步骤

1. 前置权限配置

  • 给Databricks用到的身份(托管身份/服务主体)分配两个权限:
    • Event Hub的Azure Event Hubs数据接收者角色
    • ADLS Gen2目标存储容器的存储Blob数据参与者角色
  • 提前记录好配置参数:Event Hub连接字符串、专用消费者组名(不要用默认$Default)、ADLS Gen2存储账户名、目标容器名、检查点存储路径

2. Databricks环境准备

  • Databricks Runtime 11.x及以上版本默认集成Event Hub和ADLS Gen2依赖,无需额外安装;低版本需提前安装com.microsoft.azure:azure-eventhubs-spark_2.12:2.3.22对应Jar包
  • 在Notebook中配置Event Hub连接参数,代码示例:
eh_conf = {
  "eventhubs.connectionString" : sc._jvm.org.apache.spark.eventhubs.EventHubsUtils.encrypt("你的Event Hub连接字符串"),
  "eventhubs.consumerGroup" : "你创建的流摄入专用消费者组名"
}

3. 流数据读取与格式转换

  • 用Spark Structured Streaming读取Event Hub原始流:
raw_stream_df = spark.readStream.format("eventhubs").options(**eh_conf).load()
  • 原始数据的body字段为二进制类型,按你业务要求的格式解析即可,JSON格式解析示例如下:
from pyspark.sql.functions import col, from_json
from pyspark.sql.types import StructType, StringType, TimestampType

# 替换为你实际的业务数据schema
business_schema = StructType() \
  .add("device_id", StringType()) \
  .add("event_time", TimestampType()) \
  .add("metric", StringType())

parsed_stream_df = raw_stream_df \
  .select(from_json(col("body").cast("string"), business_schema).alias("data")) \
  .select("data.*")

如果是CSV、Avro、Protobuf等其他格式,替换对应解析逻辑即可,Protobuf需额外导入对应序列化依赖

4. 写入ADLS Gen2配置

  • 支持输出Parquet、Delta、JSON等多种格式,以按天分区的Parquet格式为例,写入代码如下:
# 提前从event_time提取分区字段event_date,不需要分区可以跳过这步
parsed_stream_df = parsed_stream_df.withColumn("event_date", col("event_time").cast("date"))

write_query = parsed_stream_df.writeStream \
  .format("parquet") \
  .option("path", "abfss://目标容器名@存储账户名.dfs.core.windows.net/目标存储目录/") \
  .option("checkpointLocation", "abfss://检查点容器名@存储账户名.dfs.core.windows.net/检查点存储目录/") \
  .partitionBy("event_date") \
  .trigger(processingTime="5分钟") # 按你的延迟要求调整触发间隔,低延迟场景可以用continuous模式
  .start()

write_query.awaitTermination()
  • 配置checkpointLocation即可保证exactly-once语义,Spark会自动管理消费偏移量,任务重启后从上次断开位置继续消费,不会丢数

5. 流任务运维配置

  • 把上述代码配置为Databricks Job,设置自动重试策略即可实现自动化运行
  • 流任务的消费延迟、吞吐量等指标可以直接在Databricks的Structured Streaming UI中查看
  • 如果需要支持schema演进,建议输出格式选Delta,天然支持schema合并和变更,无需额外改造代码

内容的提问来源于stack exchange,提问作者Sai Varun Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 04:48:02