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

EMR Notebook中Spark Structured Streaming SSL连接Kafka失败

解决Spark Structured Streaming连接SSL加密Kafka无数据流入问题

问题背景

已通过confluent_kafka成功消费SSL保护的Kafka,工作代码如下:

conf = {'bootstrap.servers': 'url:port',
        'group.id': 'group1',
        'enable.auto.commit': False,
        'security.protocol': 'SSL',
        'ssl.key.location': 'file1.key',
        'ssl.ca.location': 'file2.pem',
        'ssl.certificate.location': 'file3.cert',
        'auto.offset.reset': 'earliest'
        }

consumer = Consumer(conf)
consumer.subscribe(['my_topic'])

# 能正常读取事件
msg = consumer.poll(timeout=0)

但在EMR Notebook中用Spark Structured Streaming复现时,无数据流入。当前Spark代码如下:

%%configure -f 
{
  "conf": {
    "spark.jars.packages": "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0",
    "livy.rsc.server.connect.timeout":"600s",
    "spark.pyspark.python": "python3",
    "spark.pyspark.virtualenv.enabled": "true",
    "spark.pyspark.virtualenv.type":"native",
    "spark.pyspark.virtualenv.bin.path":"/usr/bin/virtualenv"
  }
}
from pyspark.sql import SparkSession
from pyspark.sql import SQLContext
from pyspark import SparkFiles
cert_file = 'file3.cert'
pem_file = 'file2.pem'
key_file = 'file3.key'

sc.addFile(f's3://.../{cert_file}')
sc.addFile(f's3://.../{pem_file}')
sc.addFile(f's3://.../{key_file}')
spark = SparkSession\
    .builder \
    .getOrCreate()
kafka_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "url:port") \
    .option("kafka.group.id", "group1") \
    .option("enable.auto.commit", "False") \
    .option("kafka.security.protocol", "SSL") \
    .option("kafka.ssl.key.location", SparkFiles.get(key_file)) \ # SparkFiles.get()能正常获取路径
    .option("kafka.ssl.ca.location", SparkFiles.get(pem_file)) \
    .option("kafka.ssl.certificate.location", SparkFiles.get(cert_file)) \
    .option("startingOffsets", "earliest") \
    .option("subscribe", "my_topic") \
    .load()
query = kafka_df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
    .writeStream \
    .outputMode("append") \
    .format("console") \
    .start()

问题分析

Spark Structured Streaming的Kafka连接器参数需要以kafka.为前缀对应原生Kafka参数,你的代码存在以下几个关键问题:

  1. 参数前缀缺失:enable.auto.commit未加kafka.前缀,Spark无法识别该参数,应该改为kafka.enable.auto.commit
  2. 偏移量优先级问题:startingOffsets="earliest"仅在消费组无历史偏移量时生效,如果group1之前已经提交过偏移量,Spark会优先使用历史偏移量,而非从最早位置开始消费
  3. 缺少必要的等待逻辑:query.start()后没有调用query.awaitTermination(),流式查询可能在启动后立即终止,来不及拉取数据

修正后的代码

1. 调整流式读取配置

kafka_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "url:port") \
    .option("kafka.group.id", "group1") \
    .option("kafka.enable.auto.commit", "False") \  # 补全kafka.前缀
    .option("kafka.security.protocol", "SSL") \
    .option("kafka.ssl.key.location", SparkFiles.get(key_file)) \
    .option("kafka.ssl.ca.location", SparkFiles.get(pem_file)) \
    .option("kafka.ssl.certificate.location", SparkFiles.get(cert_file)) \
    .option("kafka.auto.offset.reset", "earliest") \  # 增加原生Kafka的偏移量重置参数
    .option("startingOffsets", "earliest") \
    .option("subscribe", "my_topic") \
    .load()

2. 增加查询等待逻辑

query = kafka_df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
    .writeStream \
    .outputMode("append") \
    .format("console") \
    .start()

query.awaitTermination()  # 保持查询运行,等待数据流入

额外排查步骤

如果修正后仍无数据,按以下步骤排查:

  • 检查Spark日志:查看EMR集群的Spark Driver/Executor日志,确认是否存在SSL相关错误(如证书验证失败、文件路径不存在等)
  • 重置消费组偏移量:如果group1已有历史偏移量,可通过Kafka命令行工具重置偏移量,或更换一个全新的消费组ID
  • 验证证书分发:在代码中添加打印语句,确认SparkFiles.get(key_file)返回的路径下存在文件,且所有Executor节点都能访问该文件
  • 检查网络连通性:确认EMR集群能访问Kafka的SSL端口(通常为9093),可在集群节点上执行telnet url port测试连通性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 16:20:32