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

如何连接OpenShift上部署的Confluent-Kafka并解决Python消费者无法接收消息问题

解决OpenShift Route暴露Kafka后Python消费者无法接收消息的问题

我来帮你拆解下当前遇到的问题,一步步带你排查和解决:

1. 先解决端口与协议不匹配的核心问题

你当前的Route把外部443端口转发到了Broker的9092-tcp端口,但这个端口是Kafka默认的明文(PLAINTEXT)监听端口,而你的Python消费者用security_protocol="SSL"发起SSL连接,这就导致了协议不匹配——客户端在发SSL握手请求,Broker那边却在等着明文流量,自然无法建立有效连接,消费者就会一直处于等待状态。

修复方式:

  • 先确认Kafka Broker的SSL监听端口(Confluent Kafka默认是9093,你可以在Control Center的Broker配置页面查看,或者查看部署YAML里的端口配置)。
  • 修改你的Route配置,把targetPort改成Broker的SSL端口对应的service端口,比如:
spec:
  port:
    targetPort: 9093-tcp  # 替换为实际的Broker SSL端口

2. 配置Kafka的advertised.listeners(最容易忽略的关键)

Kafka有个很重要的机制:客户端连接bootstrap server后,Broker会把自己的advertised.listeners地址返回给客户端,后续客户端会用这个地址来建立实际的消息收发连接。如果你的Broker配置里这个地址还是OpenShift内部的服务地址(比如broker:9092),消费者拿到后会尝试连接这个内部地址,肯定访问不到。

修复方式:

  • 更新Confluent Kafka的Broker配置,添加SSL类型的外部监听器,并把advertised.listeners设置为你的Route的host和443端口:
    示例配置(可以通过Confluent Operator或者直接修改Broker的环境变量/configmap来设置):
    listeners=PLAINTEXT://0.0.0.0:9092,SSL://0.0.0.0:9093
    advertised.listeners=PLAINTEXT://broker:9092,SSL://kafka-xxx.apps.yyy.zzz:443
    
  • 同时要给Broker配置对应SSL证书:因为你用的是passthrough类型的Route,SSL终止在Broker端,所以Broker需要持有与Route域名匹配的证书。如果用OpenShift集群的自动生成证书,你可以挂载router-ca或者集群的服务证书到Broker容器,然后配置以下参数:
    ssl.keystore.location=/path/to/your-keystore.jks
    ssl.keystore.password=<证书密码>
    ssl.key.password=<密钥密码>
    ssl.truststore.location=/path/to/your-truststore.jks
    ssl.truststore.password=<信任库密码>
    
    (Confluent在OpenShift部署时,通常可以通过Operator的配置面板来简化这些SSL参数的设置)

3. 完善Python消费者的SSL信任配置

你的消费者代码只指定了security_protocol="SSL",但如果Broker用的是OpenShift集群的自签名CA证书,消费者需要信任这个CA才能完成SSL握手。

修复方式:

  1. 从OpenShift集群导出CA证书:
    执行以下命令(需要有集群权限):
    oc get secret/router-ca -n openshift-ingress -o jsonpath='{.data.tls\.crt}' | base64 -d > openshift-ca.crt
    
  2. 修改Python消费者代码,添加CA证书的信任配置:
    from kafka import KafkaConsumer
    if __name__ == '__main__':
        consumer = KafkaConsumer('my_topic_which_i_see_on_control_center',
                                 bootstrap_servers=['kafka-xxx.apps.yyy.zzz:443'],
                                 api_version=(0,10,2),
                                 enable_auto_commit=True,
                                 security_protocol="SSL",
                                 ssl_cafile="/path/to/openshift-ca.crt",  # 替换为你保存证书的路径
                                 auto_commit_interval_ms=1000,
                                 auto_offset_reset="earliest",
                                 group_id='ersin_test' )
        for msg in consumer:
            print(msg)
    

4. 额外的排查验证步骤

如果做完上面的步骤还是有问题,可以试试这些方法定位:

  • 用openssl测试SSL连接:
    openssl s_client -connect kafka-xxx.apps.yyy.zzz:443
    
    如果握手成功,会显示证书详情和连接状态;如果失败,错误信息会告诉你问题出在哪(比如证书不匹配、端口错误)。
  • 查看Kafka Broker日志:检查是否有客户端连接失败的日志,比如握手错误、监听器配置错误等。
  • 确认主题和消费偏移量:在Control Center里确认my_topic_which_i_see_on_control_center确实有未消费的消息,并且你的消费者组ersin_test的偏移量没有已经追到最新(可以在消费者组页面查看偏移量状态)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 11:12:35