如何使用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项目结构,
pom.xml管理Spark和Event Hubs依赖,PySpark代码存放在src/main/python目录下
2. PySpark执行成功结果
Spark流任务正常运行,控制台输出了从Event Hubs拉取到的消息内容和事件时间
3. 常见错误示例(权限不足)
当服务主体缺少Receiver权限时,会抛出
AuthorizationFailed错误,提示无法访问Event Hubs资源
内容的提问来源于stack exchange,提问作者azuresnowflake1
相关产品推荐
相关产品推荐

