Pod上Kafka连接Java正常Python报SSL证书验证失败排查
Pod内Kafka SSL连接Java正常Python报证书验证失败问题
问题场景
- Java环境连通性验证:在Pod中使用Apache Kafka 3.2.0官方二进制包执行测试脚本,仅配置SSL协议未指定自定义信任库,执行后正常返回Topic配置、分区副本信息,确认Java环境下Kafka SSL连接无异常。
$ wget https://dlcdn.apache.org/kafka/3.2.0/kafka_2.13-3.2.0.tgz $ tar -xzf kafka_2.13-3.2.0.tgz $ cd kafka_2.13-3.2.0/bin $ echo "security.protocol=SSL" > my.properties $ ./kafka-topics.sh --describe --topic pingone_fraud --bootstrap-server <host>:9093 --command-config my.properties # output: Topic: <topic> PartitionCount: 8 ReplicationFactor: 3 Configs: cleanup.policy=delete,segment.bytes=262144000,retention.bytes=53687063712 Topic: <topic> Partition: 0 Leader: 204 Replicas: 204,203,205 Isr: 204,203,205 . . .
- Python环境测试:使用kafka-python库编写消费者代码,指定SSL协议、CA证书路径、集群地址等参数,启动时抛出SSL证书验证失败异常。
consumer = KafkaConsumer(<topic>, security_protocol='SSL', ssl_cafile='/keys/CARoot.pem', request_timeout_ms=_request_timeout_ms, connections_max_idle_ms=_connections_max_idle_ms, bootstrap_servers='<host>:9093', group_id='fraud', auto_offset_reset=_auto_offset_reset, max_poll_interval_ms=_max_poll_interval_ms, session_timeout_ms=_session_timeout_ms, max_poll_records=_max_poll_records) print(f'************* {consumer.topics()}')
异常栈如下:
Traceback (most recent call last): File "/tmp/kafka_consumer.py", line 30, in <module> max_poll_records=_max_poll_records) File "/usr/local/lib/python3.6/site-packages/kafka/consumer/group.py", line 354, in __init__ self._client = KafkaClient(metrics=self._metrics, **self.config) File "/usr/local/lib/python3.6/site-packages/kafka/client_async.py", line 240, in __init__ self.config['api_version'] = self.check_version(timeout=check_timeout) File "/usr/local/lib/python3.6/site-packages/kafka/client_async.py", line 908, in check_version version = conn.check_version(timeout=remaining, strict=strict, topics=list(self.config['bootstrap_topics_filter'])) File "/usr/local/lib/python3.6/site-packages/kafka/conn.py", line 1171, in check_version if not self.connect_blocking(timeout_at - time.time()): File "/usr/local/lib/python3.6/site-packages/kafka/conn.py", line 333, in connect_blocking self.connect() File "/usr/local/lib/python3.6/site-packages/kafka/conn.py", line 422, in connect if self._try_handshake(): File "/usr/local/lib/python3.6/site-packages/kafka/conn.py", line 501, in _try_handshake self._sock.do_handshake() File "/usr/local/lib/python3.6/ssl.py", line 1077, in do_handshake self._sslobj.do_handshake() File "/usr/local/lib/python3.6/ssl.py", line 689, in do_handshake self._sslobj.do_handshake() ssl.SSLError: [SSL: CERTIFICATE_VERIFY_FAILED] certificate verify failed (_ssl.c:852)
- 证书来源说明:Python侧使用的CA证书,是通过keytool命令从Java默认信任库
/opt/java/lib/security/cacerts中导出的别名为vault-us的根证书,用户认为与Java环境使用的证书来源完全一致。
keytool -exportcert -keystore /opt/java/lib/security/cacerts -alias vault-us -rfc -file /keys/CARoot.pem
根因分析
核心问题是两侧SSL校验的信任源、校验逻辑存在本质差异,并非证书同源导致的异常:
- 信任源范围不一致:Java测试脚本未指定自定义信任库,默认加载整个cacerts文件中包含的上百个公共根CA、私有根CA证书完成校验;而Python侧仅加载了从cacerts中导出的单个
vault-us别名证书,如果Kafka broker证书的实际签发CA不是vault-us,而是cacerts中存储的其他CA,就会出现Java校验通过、Python校验失败的情况。这是该场景下最高发的诱因。 - 证书链补全逻辑差异:Java SSL引擎支持自动从信任库中匹配、补全不完整的证书链;而Python基于OpenSSL实现的SSL模块,仅会加载
ssl_cafile指定文件内的证书,不会自动从系统/其他路径查找缺失的中间CA证书,如果只导入根证书但broker返回的证书链缺少中间CA,会直接导致链校验失败。 - 主机名校验规则差异:kafka-python默认开启SSL主机名验证,要求broker返回证书的SAN(主题备用名称)字段必须包含客户端连接时使用的域名/IP;部分旧版本Kafka Java客户端默认未强制开启主机名校验,会出现证书域名不匹配但Java连接正常的情况。
- 证书格式异常:keytool导出时如果参数错误、文件写入不完整、存在多余字符或权限不足,会导致Python无法正确解析PEM证书文件,触发校验失败。
可行解决方案
按优先级依次排查修复:
- 确认broker证书实际签发CA
执行以下命令直接从broker端口拉取完整证书链,确认根CA的实际标识:
从输出中找到证书链最上层根证书的Subject、Issuer字段,确认其是否与导出的openssl s_client -connect <host>:9093 -showcerts < /dev/nullvault-us证书一致。 - 替换为完整信任库文件
不要仅导出单个别名证书,直接将整个Java cacerts信任库导出为PEM格式供Python使用:
将Python代码中# Java cacerts默认密码为changeit,如果环境修改过请替换为实际密码 keytool -list -rfc -keystore /opt/java/lib/security/cacerts -storepass changeit > /keys/FullCA.pemssl_cafile参数值替换为/keys/FullCA.pem后重启测试。 - 校验证书文件有效性
确认/keys/FullCA.pem文件权限对Python运行用户可读,文件内所有证书均为标准PEM格式,以-----BEGIN CERTIFICATE-----开头、-----END CERTIFICATE-----结尾,无多余乱码或keytool输出的提示内容。 - 排查主机名验证问题
如果替换完整信任库后仍报错,可临时在KafkaConsumer初始化参数中添加ssl_check_hostname=False测试连通性:- 如果临时关闭后连通正常,说明broker证书SAN字段未包含你当前连接使用的
<host>地址,请优先更换为证书中已登记的域名连接,或重新为broker签发带正确SAN字段的证书。
不建议生产环境长期关闭主机名验证,会存在SSL中间人攻击风险
- 如果临时关闭后连通正常,说明broker证书SAN字段未包含你当前连接使用的
- 版本兼容修复
你当前使用的Python 3.6已停止维护,存在多个SSL模块已知bug,建议升级到Python 3.8+及kafka-python 2.0.2+稳定版本后重试。
内容的提问来源于stack exchange,提问作者Ema Il
相关产品推荐
相关产品推荐

