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参数,你的代码存在以下几个关键问题:
- 参数前缀缺失:
enable.auto.commit未加kafka.前缀,Spark无法识别该参数,应该改为kafka.enable.auto.commit - 偏移量优先级问题:
startingOffsets="earliest"仅在消费组无历史偏移量时生效,如果group1之前已经提交过偏移量,Spark会优先使用历史偏移量,而非从最早位置开始消费 - 缺少必要的等待逻辑:
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
相关产品推荐
相关产品推荐

