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

Strimzi Kafka证书轮换后Flink集群仍报SSL握手失败问题排查

问题分析与解决方案:Flink无法加载Strimzi Kafka新证书导致SSL握手失败

问题背景

部署在Kubernetes集群上的Flink应用,从Strimzi Kafka读取数据并写入云MSK。Strimzi Kafka证书过期完成轮换后,即使手动重启JobManager(JM)和TaskManager(TM),仍出现PKIX path validation failed: Path does not chain with any of the trust anchors错误,而其他基于Quarkus的微服务可正常运行。Flink已启用Checkpoint(外部PVC存储)和HA高可用机制。

核心原因与解决步骤

1. Checkpoint/HA状态残留旧SSL配置

Flink从外部PVC恢复作业状态时,会复用旧的Kafka客户端缓存信息(包括失效的证书链),导致新启动的Pod仍使用旧证书进行SSL握手。

  • 解决方法:
    • 暂停作业,备份Checkpoint存储目录后清理相关状态文件,再重新启动作业。
    • 在Kafka消费者配置中添加以下参数,强制客户端重启时重载证书:
      ssl.truststore.reload.config=true
      ssl.keystore.reload.config=true
      

2. 证书Secret挂载未生效或Flink未触发重载

尽管配置了自动密钥重载注解,但Flink容器可能未正确挂载最新Secret,或Kafka客户端默认不监听证书变化(Quarkus有内置证书重载机制,Flink无此默认行为)。

  • 解决方法:
    • 检查Pod内证书文件的更新时间,确认是否挂载了最新Secret:
      kubectl exec <flink-jm-pod> -- ls -l /path/to/mounted/certs
      
    • 确认挂载正确后,在Flink的Kafka连接器配置中明确指定证书路径并启用重载:
      ssl.truststore.location=/path/to/truststore.jks
      ssl.truststore.password=<你的密码>
      ssl.truststore.reload.config=true
      ssl.keystore.location=/path/to/keystore.jks
      ssl.keystore.password=<你的密码>
      ssl.keystore.reload.config=true
      

3. JVM默认信任库未更新

若Flink使用JVM默认信任库(cacerts),旧证书未被移除或新证书未导入,导致证书链验证失败。

  • 解决方法:
    • 检查JVM信任库中的证书列表:
      kubectl exec <flink-jm-pod> -- keytool -list -keystore $JAVA_HOME/lib/security/cacerts -storepass changeit
      
    • 将新的Strimzi CA证书导入JVM信任库,或修改Flink配置使用自定义信任库(挂载自Secret)。

4. Kafka客户端连接缓存残留

Kafka客户端可能缓存了旧的SSL会话,即使重启Pod仍复用无效连接。

  • 解决方法:
    • 在Kafka配置中缩短空闲连接超时时间,强制重新建立连接:
      connections.max.idle.ms=30000
      
    • 强制重启所有Flink Pod,必要时清理Kubernetes节点的DNS缓存。

错误日志

2025-07-15 12:45:56,309 ERROR org.apache.kafka.clients.NetworkClient [] - [Consumer clientId=k2k-consumer-0, groupId=k2k-consumer] Connection to node -1 (strimzi-kafka-bootstrap.test/100.64.117.90:9093) failed authentication due to: SSL handshake failed 
2025-07-15 12:45:56,310 ERROR org.apache.flink.connector.base.source.reader.fetcher.SplitFetcherManager [] - Received uncaught exception. java.lang.RuntimeException: SplitFetcher thread 0 received unexpected exception while polling the records at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:168) ~[flink-connector-files-1.20.1.jar:1.20.1] at
 org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.run(SplitFetcher.java:117) [flink-connector-files-1.20.1.jar:1.20.1] at 
 java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source) [?:?] at java.util.concurrent.FutureTask.run(Unknown Source) [?:?] at
 java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) [?:?] at java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) [?:?] at 
 java.lang.Thread.run(Unknown Source) [?:?] Caused by: org.apache.kafka.common.errors.SslAuthenticationException: SSL handshake failed Caused by: 
 javax.net.ssl.SSLHandshakeException: PKIX path validation failed: java.security.cert.CertPathValidatorException: Path does not chain with any of the trust anchors

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 16:57:34