OpenReplay 自托管 Kafka 集群连接指南:PLAINTEXT 与 SSL/TLS 双通道实战
2026/9/23 4:19:27 网站建设 项目流程
  • 可观测性
  • 开发工具
  • 前端
  • 后端

【免费下载链接】openreplay

Session replay, cobrowsing and product analytics you can self-host. Best for reproducing issues and iterating on your product.

项目地址:https://gitcode.com/gh_mirrors/op/openreplay
点击查看免费下载

导读

本文档是 OpenReplay 自托管方案中 Kafka 集群(位于scripts/dockerfiles/kafka)的官方连接指南,讲解如何从另一个容器宿主机上的应用连接这套 2 节点 Kafka 集群。你将掌握:两种监听通道(PLAINTEXT 明文与 SSL 加密)的端口布局与配置要点、容器内/宿主机两种场景下的命令行验证方法、控制台生产者与消费者的完整用法,以及 Java、Python、Node.js、Go 四种语言的客户端配置示例。读完本文,你可以立即在自己的业务应用或调试容器中接入该 Kafka 集群,并对 TLS 客户端的信任链配置有源码级的理解。


1. 集群拓扑与端口布局(连接前必读)

CONNECTION_GUIDE 讨论的是一套2 节点、KRaft 模式(无 ZooKeeper)、支持复制的 Kafka 集群,其基础与 TLS 两种部署分别由 docker-compose.yml 与 docker-compose-tls.yml 定义。启动方式可参考 README.md:

# 基础集群(仅 PLAINTEXT) podman-compose up -d # TLS 集群(先由 generate-certs.sh 生成证书) ./generate-certs.sh podman-compose -f docker-compose-tls.yml up -d

1.1 标准集群(PLAINTEXT)端口

对应 docker-compose.yml 的映射:

节点容器内监听宿主机映射用途
kafka-190929092PLAINTEXT 客户端监听
kafka-190939093CONTROLLER(KRaft 内部控制器通信)
kafka-290929094PLAINTEXT 客户端监听
kafka-290939095CONTROLLER

1.2 TLS 集群端口

对应 docker-compose-tls.yml 的映射:

节点容器内监听宿主机映射用途
kafka-190929092PLAINTEXT 明文通道
kafka-190939093CONTROLLER
kafka-190949094SSL 加密通道
kafka-290929095PLAINTEXT 明文通道
kafka-290939096CONTROLLER
kafka-290949097SSL 加密通道

可以看到 TLS 集群同时暴露 PLAINTEXT 与 SSL 两种监听(见KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093,SSL://:9094KAFKA_ADVERTISED_LISTENERS),这为从明文平滑迁移到加密提供了便利——你可以先让所有客户端走 9092/9095,验证无误后再切换到 9094/9097。

端口澄清:CONNECTION_GUIDE 中 "kafka-2: 9095 (mapped from internal 9092)" 与 "kafka-2: 9097 (mapped from internal 9094)" 的描述,与 README.md 的 Ports 章节一致,对应的是TLS 集群的宿主机视角;而标准集群中 kafka-2 的宿主机明文端口是 9094。连接前请先确认自己启动的是哪个 compose 文件。

1.3 集群关键事实(来自 CLUSTER_INFO.md)

  • Kafka 版本:3.9.0(固定在 3.x 大版本)
  • 模式:KRaft,无 ZooKeeper
  • 共享集群 ID:Sjg_Rr1iQbO9xpahgDbYpQ
  • 双节点复制可用;数据目录/bitnami/kafka,以非 root 用户(UID 1001)运行
  • 基础镜像为 Chainguard Wolfi,Kafka 工具链位于容器内/usr/lib/kafka/bin/

2. 连接方式总览:PLAINTEXT 与 SSL

本集群对外提供两种连接通道,CONNECTION_GUIDE 给出了一一对应的配置模板。

2.1 PLAINTEXT 连接(无加密)

端口:

  • kafka-1: 9092
  • kafka-2: 9095(由内部 9092 映射而来,即上文 TLS 集群拓扑中的宿主机视角;标准集群则为 localhost:9094)

配置(bootstrap.servers):

bootstrap.servers=kafka-1:9092,kafka-2:9092 # 从宿主机连接时: # bootstrap.servers=localhost:9092,localhost:9095

示例——列出主题:

kafka-topics.sh --list --bootstrap-server kafka-1:9092

2.2 SSL 连接(加密)

端口:

  • kafka-1: 9094
  • kafka-2: 9097(由内部 9094 映射而来)

客户端配置文件(ssl-client.properties):

bootstrap.servers=kafka-1:9094,kafka-2:9094 security.protocol=SSL ssl.truststore.location=/path/to/ca-cert.pem ssl.truststore.type=PEM ssl.endpoint.identification.algorithm=

示例——列出主题:

kafka-topics.sh --list \ --bootstrap-server kafka-1:9094 \ --command-config ssl-client.properties

ssl.endpoint.identification.algorithm=(空值)用于关闭主机名校验。本集群默认自签名证书(由 generate-certs.sh 生成),因此客户端必须禁用主机名校验才能通过握手;生产环境使用 CA 签发的证书时,应移除该行以恢复校验。


3. 从另一个容器连接

与 broker 位于同一 Podman/Docker 网络中的容器,可以直接使用服务主机名kafka-1/kafka-2通信。

3.1 启动一个客户端容器

podman run -d --name kafka-client \ --network kafka-network-tls \ -v /path/to/certs:/certs:ro \ --entrypoint /bin/sh \ your-kafka-image:latest \ -c "while true; do sleep 3600; done"

要点说明:

  • --network kafka-network-tls必须与 TLS 集群所在的网络一致(见 docker-compose-tls.yml 中的networks定义;标准集群则为kafka-network),否则无法解析kafka-1/kafka-2主机名;
  • -v /path/to/certs:/certs:rocerts/目录(内含ca-cert.pem)以只读方式挂载进客户端容器;
  • 镜像使用集群同样的 Kafka 镜像(含/usr/lib/kafka/bin/下的工具脚本),--entrypoint /bin/sh让容器保持存活以便执行调试命令。

3.2 在容器内创建 SSL 客户端配置

podman exec kafka-client sh -c 'cat > /tmp/ssl-client.properties << EOF security.protocol=SSL ssl.truststore.location=/certs/ca-cert.pem ssl.truststore.type=PEM ssl.endpoint.identification.algorithm= EOF'

注意此处ssl.truststore.location指向容器内的挂载路径/certs/ca-cert.pem

3.3 容器内连通性测试

# PLAINTEXT podman exec kafka-client kafka-topics.sh --list --bootstrap-server kafka-1:9092 # SSL podman exec kafka-client kafka-topics.sh --list \ --bootstrap-server kafka-1:9094 \ --command-config /tmp/ssl-client.properties

两条命令分别验证明文与加密通道。若想连第二个 broker,把kafka-1换成kafka-2,端口相应使用 9092/9094 即可(bootstrap.servers 中同时列出两个 broker 可获得故障转移能力)。


4. 从宿主机连接

从宿主机访问时,必须使用映射后的端口localhost(或宿主机 IP)。

4.1 PLAINTEXT

kafka-topics.sh --list --bootstrap-server localhost:9092

4.2 SSL

先创建ssl-client.properties(宿主机路径):

security.protocol=SSL ssl.truststore.location=/full/path/to/ca-cert.pem ssl.truststore.type=PEM ssl.endpoint.identification.algorithm=

然后使用:

kafka-topics.sh --list \ --bootstrap-server localhost:9094 \ --command-config ssl-client.properties

宿主机侧需要本机已安装 Kafka 命令行工具(或使用容器内工具 +podman exec代替,见 README 的管理命令章节)。注意宿主机上没有 9094 的kafka-2对称端口时,应使用 9097(TLS 集群)或 9094(标准集群)作为第二 broker。


5. 生产者 / 消费者示例

5.1 PLAINTEXT 生产者

kafka-console-producer.sh \ --bootstrap-server kafka-1:9092 \ --topic my-topic

5.2 SSL 生产者

kafka-console-producer.sh \ --bootstrap-server kafka-1:9094 \ --topic my-topic \ --producer.config ssl-client.properties

5.3 PLAINTEXT 消费者

kafka-console-consumer.sh \ --bootstrap-server kafka-1:9092 \ --topic my-topic \ --from-beginning

5.4 SSL 消费者

kafka-console-consumer.sh \ --bootstrap-server kafka-1:9094 \ --topic my-topic \ --from-beginning \ --consumer.config ssl-client.properties

SSL 场景下,生产者通过--producer.config、消费者通过--consumer.config传入同一份ssl-client.properties。若主题尚不存在,可先用 README 中的命令创建双副本主题:

podman exec kafka-1 /usr/lib/kafka/bin/kafka-topics.sh \ --create --topic my-topic \ --bootstrap-server kafka-1:9092 \ --replication-factor 2 --partitions 3

6. 应用代码接入示例

CONNECTION_GUIDE 为四种主流语言/框架给出了可直接复制的配置。

6.1 Java / Spring Boot

spring: kafka: bootstrap-servers: kafka-1:9094,kafka-2:9094 properties: security.protocol: SSL ssl.truststore.location: /path/to/ca-cert.pem ssl.truststore.type: PEM ssl.endpoint.identification.algorithm: ""

ssl.endpoint.identification.algorithm: ""对应关闭主机名校验;生产环境使用受信任 CA 时应去掉此项。

6.2 Python(kafka-python)

from kafka import KafkaProducer, KafkaConsumer # PLAINTEXT producer = KafkaProducer( bootstrap_servers=['kafka-1:9092', 'kafka-2:9092'] ) # SSL producer = KafkaProducer( bootstrap_servers=['kafka-1:9094', 'kafka-2:9094'], security_protocol='SSL', ssl_check_hostname=False, ssl_cafile='/path/to/ca-cert.pem' )

ssl_check_hostname=False是 Python 侧关闭主机名校验的方式,与属性文件中的ssl.endpoint.identification.algorithm=作用等价;ssl_cafile指向 CA 证书 PEM 文件。

6.3 Node.js(kafkajs)

const { Kafka } = require('kafkajs') // PLAINTEXT const kafka = new Kafka({ clientId: 'my-app', brokers: ['kafka-1:9092', 'kafka-2:9092'] }) // SSL const kafka = new Kafka({ clientId: 'my-app', brokers: ['kafka-1:9094', 'kafka-2:9094'], ssl: { rejectUnauthorized: false, ca: [fs.readFileSync('/path/to/ca-cert.pem', 'utf-8')] } })

kafkajs 的ssl选项直接透传给 Node.js TLS 层:ca传入 CA 证书内容,rejectUnauthorized: false关闭证书链/主机名校验(自签名证书场景必需)。

6.4 Go(sarama)

import ( "crypto/tls" "crypto/x509" "io/ioutil" "github.com/Shopify/sarama" ) // PLAINTEXT config := sarama.NewConfig() brokers := []string{"kafka-1:9092", "kafka-2:9092"} // SSL config := sarama.NewConfig() config.Net.TLS.Enable = true caCert, _ := ioutil.ReadFile("/path/to/ca-cert.pem") caCertPool := x509.NewCertPool() caCertPool.AppendCertsFromPEM(caCert) tlsConfig := &tls.Config{ RootCAs: caCertPool, InsecureSkipVerify: true, } config.Net.TLS.Config = tlsConfig brokers := []string{"kafka-1:9094", "kafka-2:9094"}

Go 侧通过标准库crypto/tls构造配置:RootCAs放入 CA 证书池,InsecureSkipVerify: true跳过服务端证书校验(与其它语言关闭主机名校验对应)。注意ioutil.ReadFile已在新版 Go 中废弃,可替换为os.ReadFile


7. 网络要求

CONNECTION_GUIDE 归纳了三种网络场景:

同一 Docker/Podman 网络(容器间):

  • 使用主机名kafka-1kafka-2
  • 客户端容器必须与 broker 位于同一网络(如kafka-network-tls,见 docker-compose-tls.yml 的networks定义)。

宿主机访问:

  • 使用localhost+ 映射端口;
  • PLAINTEXT:9092(kafka-1)、9095(kafka-2,TLS 集群视角;标准集群为 9094);
  • SSL:9094(kafka-1)、9097(kafka-2)。

外部网络访问:

  • 更新KAFKA_ADVERTISED_LISTENERS为公网 IP/主机名——这是 broker 向客户端宣告的连接地址,若与客户端实际可达地址不一致,客户端会连接失败;
  • 确保防火墙放行 9092-9097 端口段。

原理提示:Kafka 客户端先连 bootstrap 地址,再从 broker 返回的advertised listeners建立实际连接。KAFKA_ADVERTISED_LISTENERS在 docker-compose.yml 中被设为PLAINTEXT://kafka-1:9092(TLS 版追加SSL://kafka-1:9094),这正是容器内可用、而宿主机必须改用映射端口的原因。该配置最终由容器启动脚本 start-kafka.sh 写入/tmp/server.propertiesadvertised.listeners项。


8. SSL 客户端所需的文件

客户端只需要 CA 证书这一个文件:

ca-cert.pem - Certificate Authority certificate

该文件位于仓库目录certs/中,由 generate-certs.sh 生成(脚本同时产出ca-key.pemkafka-1-cert.pemkafka-1-key.pemkafka-2-cert.pemkafka-2-key.pem,并设置私钥权限 600、证书权限 644)。

为什么只需要 CA?从源码角度看:服务端证书与私钥通过KAFKA_SSL_CERT_FILE/KAFKA_SSL_KEY_FILE提供给 broker,start-kafka.sh 的setup_ssl_from_pem函数会在容器启动时自动完成PEM → PKCS12 → JKS的转换并生成 truststore;客户端只需信任签发这些证书的 CA(即ca-cert.pem),通过ssl.truststore.type=PEM直接以 PEM 形式加载即可,无需任何 keystore 转换操作。这正是该方案 "零手工 keystore 管理" 的设计目标。


9. 故障排查(Troubleshooting)

9.1 无法解析主机名(Cannot resolve hostname)

症状:报错 "DNS resolution failed for kafka-1"。

  • 确认客户端与 broker 在同一网络(podman network ls检查);
  • 改用 IP 地址替代主机名;
  • 或从宿主机使用localhost+ 映射端口。

9.2 SSL 握手失败(SSL handshake failed)

  • 核对ssl.truststore.location路径是否正确(容器内挂载路径 vs 宿主机绝对路径);
  • 确保ssl.truststore.type=PEM
  • 添加ssl.endpoint.identification.algorithm=关闭主机名校验(自签名证书场景)。

9.3 连接超时(Connection timeout)

  • 检查 broker 是否运行:podman ps | grep kafka
  • 确认端口已暴露:podman port kafka-1-tls(TLS 集群容器名为kafka-1-tls,标准集群为kafka-1);
  • 查看 broker 日志:podman logs kafka-1-tls

其它可用诊断手段(来自 CLUSTER_INFO.md 与 README.md):

# 查看 broker 协议版本,确认可连通 podman exec kafka-1 /usr/lib/kafka/bin/kafka-broker-api-versions.sh \ --bootstrap-server kafka-1:9092 # 检查 SSL 端口是否在监听 podman exec kafka-1-tls netstat -tlnp | grep 9094

10. 总结:最小 SSL 配置速查

SSL 客户端最小配置:

security.protocol=SSL ssl.truststore.location=/path/to/ca-cert.pem ssl.truststore.type=PEM ssl.endpoint.identification.algorithm=

网络地址速查:

  • 与 broker 同网络:kafka-1:9094kafka-2:9094(容器间);
  • 宿主机:localhost:9094localhost:9097(TLS 集群映射端口)。

连接要点回顾:

  1. 明文走 9092/9095(容器间)或 9092/9094(宿主机标准集群),加密走 9094/9097;
  2. 客户端只需分发ca-cert.pem一个文件,配合ssl.truststore.type=PEM免去 keystore 转换;
  3. 自签名证书下务必关闭主机名校验(ssl.endpoint.identification.algorithm=或各语言等价选项),生产环境换成 CA 签发证书后应恢复校验;
  4. 跨网络/外部访问必须同步修改KAFKA_ADVERTISED_LISTENERS并放行 9092-9097 端口。

更多背景信息可继续阅读同目录下的 TLS_SETUP.md(TLS 完整配置与生产检查清单)、INSECURE_TLS.md(开发环境免证书校验方案)、CUSTOM_CONFIG.md(消息大小、留存策略等自定义配置)以及 CLUSTER_INFO.md(集群状态与管理命令)。

  • 可观测性
  • 开发工具
  • 前端
  • 后端

【免费下载链接】openreplay

Session replay, cobrowsing and product analytics you can self-host. Best for reproducing issues and iterating on your product.

项目地址:https://gitcode.com/gh_mirrors/op/openreplay
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询