使用服务主体+证书认证的Spark结构化流作业消费EventHub消息
在Azure Synapse中用PySpark结构化流+AAD认证消费EventHub的实现方案
1. 依赖准备
确保你的Synapse Spark池已安装以下依赖:
azure-eventhubs-spark:版本需匹配Spark池的Spark版本(Spark 3.x建议用2.3.22及以上兼容版本)azure-identity:用于AAD token获取
可在Spark池的「包」选项卡通过PyPI安装,或在Notebook开头执行:
%pip install azure-eventhubs-spark azure-identity
2. 编写AAD Token回调函数
创建符合EventHub Spark库要求的令牌获取回调类:
from azure.identity import DefaultAzureCredential class EventHubAADTokenProvider: def __init__(self, eventhub_namespace): self.namespace = eventhub_namespace self.credential = DefaultAzureCredential() # 必须保持该方法签名:接收两个参数(此处未用到),返回token字符串 def getToken(self, _, __): # EventHub的AAD权限范围 scope = "https://eventhubs.azure.net/.default" token_response = self.credential.get_token(scope) return token_response.token
3. 打包回调函数并加载到Spark池
如需将回调函数作为复用包加载:
- 按如下目录结构组织代码:
eventhub_aad_utils/ ├── __init__.py └── token_provider.py # 存放上述EventHubAADTokenProvider类 - 将整个目录打包成
eventhub_aad_utils.zip - 在Synapse Spark池的「包」选项卡上传该zip包,或在Notebook中通过
spark.sparkContext.addPyFile("path/to/eventhub_aad_utils.zip")加载
4. 结构化流消费EventHub的完整代码
from pyspark.sql import SparkSession from eventhub_aad_utils.token_provider import EventHubAADTokenProvider # 初始化Spark会话(Synapse中可省略,默认已提供) spark = SparkSession.builder.appName("EH-Structured-Streaming").getOrCreate() # 配置EventHub参数 EVENTHUB_NAMESPACE = "your-eventhub-namespace" EVENTHUB_NAME = "your-eventhub-name" # 创建token提供器实例 token_provider = EventHubAADTokenProvider(EVENTHUB_NAMESPACE) # 构建EventHub流读取选项 eh_stream_options = { "eventhubs.namespace": EVENTHUB_NAMESPACE, "eventhubs.name": EVENTHUB_NAME, "eventhubs.auth.type": "OAuth", "eventhubs.auth.token.provider": token_provider, # 设置起始位置:从最新消息开始消费 "eventhubs.startingPosition": """{"offset":"-1","seqNo":-1,"enqueuedTime":null,"isInclusive":true}""" } # 读取EventHub流 raw_stream_df = spark.readStream.format("eventhubs").options(**eh_stream_options).load() # 解析消息体(示例:假设消息是JSON格式,替换为你的实际Schema) from pyspark.sql.functions import col, from_json message_schema = "id INT, content STRING, timestamp TIMESTAMP" parsed_stream_df = raw_stream_df.select( from_json(col("body").cast("string"), message_schema).alias("message") ).select("message.*") # 输出到控制台(测试用,生产环境可替换为ADLS/SQL池等) stream_query = parsed_stream_df.writeStream \ .outputMode("append") \ .format("console") \ .trigger(processingTime="5 seconds") \ .start() stream_query.awaitTermination()
关键注意事项
- 确保Spark池使用的托管标识/服务主体拥有EventHub的Azure Event Hubs Data Receiver权限
DefaultAzureCredential会自动使用Synapse Spark池的托管标识,无需手动配置密钥- 如果需要从特定偏移量消费,调整
eventhubs.startingPosition参数即可
内容的提问来源于stack exchange,提问作者s528060
相关产品推荐
相关产品推荐

