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

使用服务主体+证书认证的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池

如需将回调函数作为复用包加载:

  1. 按如下目录结构组织代码:
    eventhub_aad_utils/
    ├── __init__.py
    └── token_provider.py  # 存放上述EventHubAADTokenProvider类
    
  2. 将整个目录打包成eventhub_aad_utils.zip
  3. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 02:48:22