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

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校验的信任源、校验逻辑存在本质差异,并非证书同源导致的异常:

  1. 信任源范围不一致:Java测试脚本未指定自定义信任库,默认加载整个cacerts文件中包含的上百个公共根CA、私有根CA证书完成校验;而Python侧仅加载了从cacerts中导出的单个vault-us别名证书,如果Kafka broker证书的实际签发CA不是vault-us,而是cacerts中存储的其他CA,就会出现Java校验通过、Python校验失败的情况。这是该场景下最高发的诱因。
  2. 证书链补全逻辑差异:Java SSL引擎支持自动从信任库中匹配、补全不完整的证书链;而Python基于OpenSSL实现的SSL模块,仅会加载ssl_cafile指定文件内的证书,不会自动从系统/其他路径查找缺失的中间CA证书,如果只导入根证书但broker返回的证书链缺少中间CA,会直接导致链校验失败。
  3. 主机名校验规则差异:kafka-python默认开启SSL主机名验证,要求broker返回证书的SAN(主题备用名称)字段必须包含客户端连接时使用的域名/IP;部分旧版本Kafka Java客户端默认未强制开启主机名校验,会出现证书域名不匹配但Java连接正常的情况。
  4. 证书格式异常:keytool导出时如果参数错误、文件写入不完整、存在多余字符或权限不足,会导致Python无法正确解析PEM证书文件,触发校验失败。

可行解决方案

按优先级依次排查修复:

  1. 确认broker证书实际签发CA
    执行以下命令直接从broker端口拉取完整证书链,确认根CA的实际标识:
    openssl s_client -connect <host>:9093 -showcerts < /dev/null
    
    从输出中找到证书链最上层根证书的Subject、Issuer字段,确认其是否与导出的vault-us证书一致。
  2. 替换为完整信任库文件
    不要仅导出单个别名证书,直接将整个Java cacerts信任库导出为PEM格式供Python使用:
    # Java cacerts默认密码为changeit,如果环境修改过请替换为实际密码
    keytool -list -rfc -keystore /opt/java/lib/security/cacerts -storepass changeit > /keys/FullCA.pem
    
    将Python代码中ssl_cafile参数值替换为/keys/FullCA.pem后重启测试。
  3. 校验证书文件有效性
    确认/keys/FullCA.pem文件权限对Python运行用户可读,文件内所有证书均为标准PEM格式,以-----BEGIN CERTIFICATE-----开头、-----END CERTIFICATE-----结尾,无多余乱码或keytool输出的提示内容。
  4. 排查主机名验证问题
    如果替换完整信任库后仍报错,可临时在KafkaConsumer初始化参数中添加ssl_check_hostname=False测试连通性:
    • 如果临时关闭后连通正常,说明broker证书SAN字段未包含你当前连接使用的<host>地址,请优先更换为证书中已登记的域名连接,或重新为broker签发带正确SAN字段的证书。
      不建议生产环境长期关闭主机名验证,会存在SSL中间人攻击风险
  5. 版本兼容修复
    你当前使用的Python 3.6已停止维护,存在多个SSL模块已知bug,建议升级到Python 3.8+及kafka-python 2.0.2+稳定版本后重试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 20:48:28