使用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
相关产品推荐
相关产品推荐

