Databricks中PySpark通过Kafka插件读取Heroku PostgreSQL流数据异常排查
问题分析与修复方案
1. SSL配置参数错误(核心问题)
你当前使用的Keystore相关参数是针对JKS格式证书的,但你的证书是PEM格式,参数完全不匹配,导致Spark无法建立与Kafka的SSL连接,因此一直处于初始化状态,无法获取数据。正确的PEM格式SSL配置应替换所有keystore相关参数,改用以下配置:
df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", kafkaBrokers) \ .option("subscribe", topic) \ .option("kafka.group.id", topic) \ .option("includeHeaders", "true") \ .option("startingOffsets", "earliest") \ .option("kafka.ssl.protocol", "SSL") \ .option("kafka.ssl.enabled.protocols", "SSL") \ .option("kafka.ssl.endpoint.identification.algorithm", "") \ # PEM格式CA证书配置 .option("kafka.ssl.truststore.type", "PEM") \ .option("kafka.ssl.truststore.location", "/databricks/ca.pem") \ # PEM格式客户端证书与私钥配置 .option("kafka.ssl.certificate.location", "/databricks/cert.pem") \ .option("kafka.ssl.key.location", "/databricks/key.pem") \ # 如果私钥设置了密码,添加此行,否则删除 # .option("kafka.ssl.key.password", "你的私钥密码") \ .load()
注意:
- 移除所有
kafka.ssl.truststore、kafka.ssl.keystore、kafka.ssl.keystore.key这类JKS容器专属参数。 - 用
dbutils.fs.ls("/databricks/")验证三个PEM文件是否存在于指定路径。
2. Trigger相关问题
.trigger(continuous="1 second")属于WriteStream的内置方法,无需额外导入,直接链式调用即可。但注意:
- Continuous触发仅支持Spark 2.3+、Kafka 0.10.2+版本。
- 当前核心问题是连接失败,先解决SSL配置再考虑触发规则,正确写法如下:
q=df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \ .writeStream \ .format("console") \ .outputMode("append") \ .trigger(continuous="1 second") \ .start()
3. 其他排查点
- 确认
kafkaBrokers和topic变量值与kafka-python使用的配置完全一致。 - 用Kafka命令行工具验证Topic存在且有消息:
kafka-topics.sh --list --bootstrap-server <brokers> --command-config <ssl-config-file>。 - 查看Databricks作业的驱动日志,检查是否有SSL连接失败的具体报错(如证书不匹配、文件缺失等)。
验证步骤
- 替换SSL配置后重新运行代码。
- 观察
q.status输出,若isDataAvailable变为true,说明连接成功并开始接收数据。 - 控制台将逐步输出流数据。
内容的提问来源于stack exchange,提问作者Shaggy
相关产品推荐
相关产品推荐

