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

使用PySpark从Azure Event Hub读取数据无输出问题排查

PySpark读取Azure Event Hub无数据问题排查

你的代码如下:

EH_CONN_STR = 'Endpoint=sb://event-hub-18-jul.servicebus.windows.net/;SharedAccessKeyName=eh_policy_18_july;SharedAccessKey=TOc+O/+U+QuuZ5R33HsiwUjsc1C8qRhCy+AEhFxkLRE=;EntityPath=ehub-18'
EH_NAMESPACE = 'event-hub-18-jul'

EH_NAME ='ehub-18'

KAFKA_OPTIONS = {
  "kafka.bootstrap.servers"  : f"{EH_NAMESPACE}.servicebus.windows.net:9093",
  "subscribe"                : EH_NAME,
  "kafka.sasl.mechanism"     : "PLAIN",
  "kafka.security.protocol"  : "SASL_SSL",
  "kafka.sasl.jaas.config"   : f"kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"{EH_CONN_STR}\";",
  "kafka.request.timeout.ms": "60000",
  "kafka.session.timeout.ms": "30000",
  "kafka.metadata.max.age.ms": "10000"
}
df = spark.read.format("kafka").options(**KAFKA_OPTIONS).load()
display(df)

针对代码持续运行无数据的问题,可按以下步骤排查:

1. 修正JAAS配置中的登录模块类名

Spark的Kafka集成通常不需要kafkashaded前缀,将kafka.sasl.jaas.config中的类名改为org.apache.kafka.common.security.plain.PlainLoginModule,避免类加载失败导致的连接问题。

2. 指定批处理读取的起始偏移量

当前使用的spark.read是批处理模式,默认从最新偏移量开始读取。如果Event Hub中的历史消息已被消费(或偏移量指向队列末尾),代码会一直等待新消息。添加kafka.startingOffsets参数指定读取起始位置:

  • 设置为"earliest":读取所有历史消息
  • 设置为"latest":仅读取新产生的消息

3. 区分批处理与流式读取场景

  • 如果是一次性读取现有数据:保持spark.read,但需确保起始偏移量配置正确
  • 如果是持续消费实时数据:改用spark.readStream流式读取,避免批处理模式下的无限等待

4. 检查依赖兼容性

确保spark-sql-kafka-0-10依赖版本与Spark版本匹配,版本不兼容可能导致连接或认证失败。

5. 查看Spark日志定位错误

检查Spark Driver和Executor日志,排查是否存在认证失败、连接超时、Topic不存在等具体错误信息,这些日志是定位问题的关键。

修改后的示例代码(批处理模式)

EH_CONN_STR = 'Endpoint=sb://event-hub-18-jul.servicebus.windows.net/;SharedAccessKeyName=eh_policy_18_july;SharedAccessKey=TOc+O/+U+QuuZ5R33HsiwUjsc1C8qRhCy+AEhFxkLRE=;EntityPath=ehub-18'
EH_NAMESPACE = 'event-hub-18-jul'
EH_NAME ='ehub-18'

KAFKA_OPTIONS = {
  "kafka.bootstrap.servers": f"{EH_NAMESPACE}.servicebus.windows.net:9093",
  "subscribe": EH_NAME,
  "kafka.sasl.mechanism": "PLAIN",
  "kafka.security.protocol": "SASL_SSL",
  "kafka.sasl.jaas.config": f"org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"{EH_CONN_STR}\";",
  "kafka.request.timeout.ms": "60000",
  "kafka.session.timeout.ms": "30000",
  "kafka.metadata.max.age.ms": "10000",
  "kafka.startingOffsets": "earliest"  # 读取所有历史消息
}

df = spark.read.format("kafka").options(**KAFKA_OPTIONS).load()
df.printSchema()  # 先查看数据结构
df.show(10, truncate=False)  # 显示前10条数据

流式读取示例(持续消费)

# 配置同上述KAFKA_OPTIONS
stream_df = spark.readStream.format("kafka").options(**KAFKA_OPTIONS).load()
# 将消息输出到控制台
query = stream_df.writeStream.outputMode("append").format("console").start()
query.awaitTermination()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 20:55:57