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

如何使用Azure Databricks(PySpark)通过Azure AD服务主体消费Event Hubs

PySpark 通过 Azure AD 服务主体消费 Azure Event Hubs

代码实现

以下是纯Python的PySpark代码,通过Azure AD服务主体认证消费Event Hubs事件:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

# 初始化SparkSession,引入对应版本的Event Hubs连接器
spark = SparkSession.builder \
    .appName("EH-AAD-Consumer") \
    .config("spark.jars.packages", "com.microsoft.azure:azure-eventhubs-spark_2.12:2.3.22") \
    # 本地运行可添加该配置:.config("spark.master", "local[*]")
    .getOrCreate()

# 替换为你的Azure AD服务主体和Event Hubs信息
AAD_TENANT_ID = "your-tenant-id"
AAD_CLIENT_ID = "your-service-principal-client-id"
AAD_CLIENT_SECRET = "your-service-principal-secret"
EH_NAMESPACE = "your-eventhubs-namespace"
EH_NAME = "your-eventhub-name"

# Event Hubs连接配置
eh_config = {
    "eventhubs.namespace": EH_NAMESPACE,
    "eventhubs.name": EH_NAME,
    "eventhubs.auth.method": "AAD",
    "eventhubs.auth.aad.tenant.id": AAD_TENANT_ID,
    "eventhubs.auth.aad.client.id": AAD_CLIENT_ID,
    "eventhubs.auth.aad.client.secret": AAD_CLIENT_SECRET,
    # 从最新位置开始消费,可根据需求调整startingPosition
    "eventhubs.startingPosition": "{\"offset\":\"-1\",\"seqNo\":-1,\"enqueuedTime\":null,\"isInclusive\":true}"
}

# 读取Event Hubs流数据
raw_stream_df = spark.readStream \
    .format("eventhubs") \
    .options(**eh_config) \
    .load()

# 解析消息体(根据实际消息格式调整,这里假设为字符串格式)
parsed_stream_df = raw_stream_df.select(
    col("body").cast("string").alias("raw_message"),
    col("enqueuedTime").alias("event_time")
)

# 输出到控制台(生产环境可替换为ADLS、Kafka等sink)
console_query = parsed_stream_df.writeStream \
    .outputMode("append") \
    .format("console") \
    .option("truncate", "false") \
    .start()

console_query.awaitTermination()

关键注意事项

  • 确保spark.jars.packages中的连接器版本与你的Spark版本兼容(示例适配Spark 3.x)
  • 服务主体必须拥有Event Hubs的Azure Event Hubs Data Receiver权限,可通过Azure门户在Event Hubs或命名空间级别分配
  • 集群环境需确保能正常访问Azure AD和Event Hubs服务

相关截图

1. Java项目结构(Maven示例)

Maven项目结构截图

标准Maven项目结构,pom.xml管理Spark和Event Hubs依赖,PySpark代码存放在src/main/python目录下

2. PySpark执行成功结果

PySpark消费事件成功截图

Spark流任务正常运行,控制台输出了从Event Hubs拉取到的消息内容和事件时间

3. 常见错误示例(权限不足)

权限错误提示截图

当服务主体缺少Receiver权限时,会抛出AuthorizationFailed错误,提示无法访问Event Hubs资源

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 19:07:56