OpenReplay 自托管 Kafka Helm Chart 实战:基于 StatefulSet 的 KRaft 模式与可选 TLS 部署指南
【免费下载链接】openreplaySession 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 承担着会话录制数据、事件流与产品分析消息的传输通道,是整个数据流水线的核心消息总线。本篇文章基于 OpenReplay 仓库内 scripts/helmcharts/databases/charts/kafka/IMPLEMENTATION_SUMMARY.md 的实现总结,结合 Chart 模板、values.yaml 与 参考清单 等仓库资源,完整讲解该 Kafka Helm Chart 从样板代码演进为可生产使用的 StatefulSet 部署的全过程。读完本文,你将掌握:KRaft 模式(免 ZooKeeper)下的 Kafka 集群编排原理、三类 Service 的拓扑设计、控制器仲裁投票(Controller Quorum Voters)的动态生成机制、可选 TLS 加密的证书注入流程,以及如何将该 Chart 集成进 databases 父 Chart 并完成生产级部署。
一、设计背景:从样板 Chart 到生产级 Kafka 部署
在改造之前,charts/kafka目录下的 Chart 只是 Helm 创建命令生成的通用样板:一个无状态 Deployment、一个 Ingress、一个 HPA,外加一堆占位测试文件——这些结构完全不适合 Kafka 这种有状态、需要稳定网络标识与持久化存储的消息系统。
本次实现的目标非常明确,依据 IMPLEMENTATION_SUMMARY.md:
- 将 Deployment 替换为StatefulSet,为每个 Broker 提供稳定的 Pod 身份(
kafka-0、kafka-1…)与专属 PVC; - 全面支持KRaft 模式,不再依赖 ZooKeeper;
- 提供可选的 TLS 加密能力,并向下兼容 PLAINTEXT;
- 所有配置项参数化到 values.yaml,做到开箱即用、可定制。
改造的实现基准来自仓库内的两份原生 K8s 参考清单:
- scripts/dockerfiles/kafka/kube/k8s-kafka-kraft.yaml:基础 KRaft 配置;
- scripts/dockerfiles/kafka/kube/k8s-kafka-kraft-tls.yaml:TLS 增强配置。
Chart 元信息也在 Chart.yaml 中同步更新:version: 11.8.6、appVersion: "3"(对应 Kafka 3.x 镜像),描述为 "Apache Kafka StatefulSet with KRaft mode support"。
文件变动一览(依据实现摘要):
| 类别 | 文件 |
|---|---|
| 新建 | templates/statefulset.yaml、templates/service-headless.yaml、templates/service-ssl.yaml、values-tls.yaml、README.md、IMPLEMENTATION_SUMMARY.md |
| 修改 | templates/_helpers.tpl(新增自定义函数)、templates/service.yaml(更新 Kafka 端口)、templates/serviceaccount.yaml(增加组件标签)、templates/NOTES.txt(Kafka 专属信息)、values.yaml(全面重写)、Chart.yaml |
| 移除 | templates/deployment.yaml(由 StatefulSet 取代)、templates/ingress.yaml(Kafka 无需 Ingress)、templates/hpa.yaml(不适用于 StatefulSet)、templates/httproute.yaml、templates/tests/ |
二、整体拓扑:三类 Service + 一个 StatefulSet
该 Chart 的运行时拓扑由 1 个 StatefulSet 和 3 个 Service 构成,各司其职:
- Headless Service(
<release>-headless):clusterIP: None,为 StatefulSet 提供稳定的 DNS 记录(kafka-0.<release>-headless.<namespace>.svc.cluster.local),是 Broker 间通信与客户端直连 Pod 的基础; - Client Service(
<release>):标准 ClusterIP,暴露 9092 端口(PLAINTEXT),供集群内客户端(如 OpenReplay 的 API、Assets、Assist 等组件)统一接入; - SSL Service(
<release>-ssl):仅在启用 TLS 时渲染,暴露 9094 端口(SSL),为加密客户端提供专用入口。
部署后从集群内访问的端点分别为(以db命名空间为例):
kafka.db.svc.cluster.local:9092 # 客户端 PLAINTEXT kafka-headless.db.svc.cluster.local:9092 # Headless(直连 Pod) kafka-ssl.db.svc.cluster.local:9094 # SSL(启用 TLS 时) kafka-0.kafka-headless.db.svc.cluster.local:9093 # 单 Pod(INTERNAL 控制器监听)这些访问方式与 README.md 中记录的端点完全一致,也保留了原始参考清单 k8s-kafka-kraft.yaml 的命名约定(kafka-headless、kafka、kafka-ssl)。
三、核心组件逐项拆解
3.1 StatefulSet:KRaft 模式的载体
statefulset.yaml 是整个 Chart 的心脏,关键设计如下:
podManagementPolicy: Parallel:所有 Pod 并行创建,加快集群启动速度(默认值来自 values.yaml);serviceName指向 Headless Service,这是 StatefulSet 获得稳定 DNS 身份的前提;updateStrategy: RollingUpdate:滚动升级 Broker;replicas可配置(默认 2),通过--set replicaCount=3即可扩容;- Pod 反亲和(anti-affinity):
preferredDuringSchedulingIgnoredDuringExecution,按kubernetes.io/hostname拓扑键将不同 Kafka Pod 尽量调度到不同节点,提升高可用性; volumeClaimTemplates:每个副本自动申请独立 PVC(默认 100Gi、ReadWriteOnce),数据目录挂载到/bitnami/kafka。
容器环境变量的设计直接继承了参考清单:KAFKA_NODE_ID取自 Pod 名(保证kafka-0的节点 ID 恒为 0),MY_POD_NAME、MY_POD_IP通过 downward API 注入。KRaft 相关变量在kraft.enabled=true时注入:
- name: KAFKA_CLUSTER_ID value: "Sjg_Rr1iQbO9xpahgDbYpQ" - name: KAFKA_PROCESS_ROLES value: "broker,controller" # 单进程同时承担 broker 与 controller 角色 - name: KAFKA_CONTROLLER_QUORUM_VOTERS value: {{ include "kafka.controllerQuorumVoters" . | quote }} # 动态生成 - name: KAFKA_CONTROLLER_LISTENER_NAMES value: "INTERNAL"其中KAFKA_CLUSTER_ID必须是合法的 Base64 UUID,可用kafka-storage.sh random-uuid生成;processRoles采用broker,controller组合模式,即每个 Broker 同时也是控制器投票者,无需单独的控制器节点。
存储、探针与资源限制同样参数化:日志目录为/bitnami/kafka/logs(emptyDir 卷),livenessProbe与readinessProbe均以 TCP 探测kafka-client端口(9092),默认参数为 initialDelay 30s/20s、period 10s、failureThreshold 3/6;资源请求默认 500m CPU / 1Gi 内存,上限 2000m CPU / 2Gi 内存。
3.2 Headless Service:Broker 发现的基石
service-headless.yaml 的关键点在于publishNotReadyAddresses: true。在 KRaft 模式下,第一个 Broker 启动时会作为控制器参与选主,此时 Pod 尚未就绪(readiness 未通过),但其他 Broker 必须能通过 DNS 找到它来完成集群引导——该开关正是为此而设。
该 Service 按需暴露全部三个端口(tcp-client: 9092、tcp-internal: 9093、tcp-ssl: 9094),targetPort引用容器端口的命名引用(kafka-client/kafka-internal/kafka-ssl),与参考清单 k8s-kafka-kraft-tls.yaml 中 Headless Service 的端口设计一致。
3.3 Client / SSL Service
- service.yaml:类型默认
ClusterIP(可用service.type覆盖),端口取service.client.port(9092),sessionAffinity: None; - service-ssl.yaml:整体被
{{- if .Values.listeners.ssl.enabled }}包裹,未启用 TLS 时不会产生任何资源。
两个 Service 的 selector 均同时匹配selectorLabels与componentLabels,确保与 StatefulSet Pod 精确关联。
3.4 ServiceAccount 与标签体系
serviceaccount.yaml 支持serviceAccount.create开关、automountServiceAccountToken与 annotations,并统一追加app.kubernetes.io/component: kafka组件标签。这套标签体系贯穿所有资源,也是 QUICKSTART.md 中kubectl get pods -n db -l app.kubernetes.io/name=kafka一类过滤命令得以生效的基础。
3.5 Helm 辅助函数:动态生成的四个关键模板
_helpers.tpl 在标准 name/fullname/labels 函数之外,新增了四个与 Kafka 强相关的模板函数:
kafka.componentLabels:输出app.kubernetes.io/component: kafka;kafka.controllerQuorumVoters:动态生成仲裁投票者列表。对每个副本 i(0 起),生成节点ID@全限定DNS:9093条目,节点 ID 为i+1,DNS 为<fullname>-<i>.<fullname>-headless.<namespace>.svc.cluster.local,最终用逗号连接。默认 2 副本时输出:1@kafka-0.kafka-headless.db.svc.cluster.local:9093,2@kafka-1.kafka-headless.db.svc.cluster.local:9093这与参考清单中的硬编码值完全等价,但随
replicaCount自动伸缩;kafka.advertisedListeners:将CLIENT://${MY_POD_NAME}.<fullname>-headless.<namespace>.svc.cluster.local:9092(及启用 TLS 时的SSL://...:9094)组装为通告地址,MY_POD_NAME在容器内展开,确保客户端拿到的每个 Broker 地址都指向正确的 Pod;kafka.listeners/kafka.listenerSecurityProtocolMap:按listeners.client/internal/ssl.enabled开关组装KAFKA_LISTENERS(CLIENT://:9092,INTERNAL://:9093,SSL://:9094)与协议映射(CLIENT:PLAINTEXT,INTERNAL:PLAINTEXT,SSL:SSL)。
值得注意的是,StatefulSet 中KAFKA_INTER_BROKER_LISTENER_NAME由listeners.client.enabled决定:客户端监听启用时使用CLIENT,否则回退到INTERNAL——与参考清单中基础版用CLIENT、TLS 版用INTERNAL的做法保持了行为兼容。
四、values.yaml 全量配置参考
values.yaml 是配置的唯一事实来源,主要分段如下:
4.1 集群与镜像
| 参数 | 默认值 | 说明 |
|---|---|---|
replicaCount | 2 | Broker 副本数,生产建议 3+ |
image.repository | ghcr.io/openreplay/kafka | Kafka 镜像仓库 |
image.tag | "3" | Kafka 主版本 |
image.pullPolicy | IfNotPresent | 拉取策略 |
serviceAccount.create | true | 是否创建 ServiceAccount |
4.2 KRaft 配置
kraft: enabled: true clusterId: "Sjg_Rr1iQbO9xpahgDbYpQ" # 必须为合法 Base64 UUID processRoles: "broker,controller" controllerListenerNames: "INTERNAL"4.3 Listener 配置
| 参数 | 默认值 | 说明 |
|---|---|---|
listeners.client.enabled/port/protocol | true/9092/PLAINTEXT | 客户端监听 |
listeners.internal.enabled/port/protocol | true/9093/PLAINTEXT | 内部(broker↔broker、控制器)监听 |
listeners.ssl.enabled/port/protocol | false/9094/SSL | 可选 SSL 监听 |
4.4 TLS 配置
tls: enabled: false secretName: kafka-tls-certs clientAuth: "required" # 客户端证书认证策略 endpointIdentificationAlgorithm: "" # 空字符串表示关闭主机名校验Secret 中期望的键为:ca-cert.pem、kafka-{0,1,2,...}-cert.pem、kafka-{0,1,2,...}-key.pem,即每个副本需持有独立证书。
4.5 Kafka 服务端参数(kafka:段)
该段与参考清单 k8s-kafka-kraft.yaml 中的环境变量一一对应,均为字符串形式的引号值:
| 分组 | 参数 | 默认值 | 说明 |
|---|---|---|---|
| 消息大小 | messageMaxBytes/replicaFetchMaxBytes | 3145728(3MB) | 单条消息与副本拉取上限 |
| 保留策略 | logRetentionHours/logRetentionBytes/logSegmentBytes | 168(7 天)/1073741824(1GB)/1073741824 | 时间与大小双维度保留 |
| 刷盘 | logFlushIntervalMessages/logFlushIntervalMs/logRetentionCheckIntervalMs | 10000/1000/300000 | 消息数与毫秒触发刷盘 |
| 副本因子 | defaultReplicationFactor/offsetsTopicReplicationFactor/transactionStateLogReplicationFactor/transactionStateLogMinIsr | 均为1 | 2 节点集群采用单副本 |
| 性能 | numIoThreads/numNetworkThreads/numPartitions/numRecoveryThreadsPerDataDir | 8/3/1/1 | 线程与分区数 |
| 网络缓冲 | socketReceiveBufferBytes/socketRequestMaxBytes/socketSendBufferBytes | 102400/104857600/102400 | 缓冲区字节数 |
| 安全 | autoCreateTopicsEnable/deleteTopicEnable/allowEveryoneIfNoAclFound/superUsers | true/false/true/User:admin | ACL 与超级用户 |
这些参数在 StatefulSet 中被渲染为KAFKA_*或KAFKA_CFG_*环境变量(如KAFKA_MESSAGE_MAX_BYTES、KAFKA_CFG_NUM_IO_THREADS),最终由镜像内的启动脚本转换为 Kafka server 配置。
4.6 资源、持久化、亲和与探针
resources.requests:500m CPU / 1Gi 内存;resources.limits:2000m CPU / 2Gi 内存;persistence:enabled: true、storageClass: ""(空则使用集群默认 StorageClass)、accessModes: [ReadWriteOnce]、size: 100Gi、支持annotations;affinity.podAntiAffinity:preferredDuringSchedulingIgnoredDuringExecution,topologyKey 为kubernetes.io/hostname,权重 1;livenessProbe/readinessProbe:TCP 探测kafka-client端口;podManagementPolicy: Parallel、updateStrategy.type: RollingUpdate;- 扩展钩子:
extraEnvVars、extraVolumes、extraVolumeMounts可注入自定义配置;nodeSelector、tolerations支持调度约束。
五、可选 TLS:从证书注入到加密通信
TLS 能力完全复刻并参数化了参考清单 k8s-kafka-kraft-tls.yaml 的实现。
5.1 证书注入的 init container 原理
当tls.enabled=true时,StatefulSet 会额外渲染一个setup-certsinit 容器(statefulset.yaml 第 47–72 行):
POD_ID=$(echo $POD_NAME | sed 's/<fullname>-//') # 提取 Pod 序号 cp /tls-secret/ca-cert.pem /tls/ca-cert.pem cp /tls-secret/kafka-${POD_ID}-cert.pem /tls/server-cert.pem cp /tls-secret/kafka-${POD_ID}-key.pem /tls/server-key.pem chmod 644 /tls/*.pem chmod 600 /tls/server-key.pem即:从名为tls.secretName的 Secret 中,按 Pod 序号挑出对应的kafka-<N>-cert.pem/kafka-<N>-key.pem,连同 CA 证书一起拷贝到共享的 emptyDir 卷/tls,主容器以只读方式挂载该卷。Secret 卷以defaultMode: 0644挂载。
随后注入 TLS 环境变量:
KAFKA_SSL_CERT_FILE=/tls/server-cert.pem KAFKA_SSL_KEY_FILE=/tls/server-key.pem KAFKA_SSL_CA_FILE=/tls/ca-cert.pem KAFKA_SSL_CLIENT_AUTH=required # 由 tls.clientAuth 控制 KAFKA_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM= # 空 = 关闭主机名校验从镜像侧看,scripts/dockerfiles/kafka/Dockerfile 基于 Wolfi 基础镜像安装 Kafka 3、OpenSSL 与 tini,入口脚本 start-kafka.sh 会检测到 PEM 证书后自动执行openssl pkcs12 -export+keytool -importkeystore,将 PEM 转换为 JKS keystore/truststore(默认密码kafka-ssl-pass,可用KAFKA_SSL_KEYSTORE_PASSWORD覆盖),从而让 Kafka 原生支持 PEM 输入。
5.2 证书生成与 Secret 创建
可复用仓库自带的 generate-certs.sh:脚本生成有效期 365 天的 CA 证书,并为kafka-1、kafka-2(即副本 0/1)签发带 SAN(DNS:kafka-N, DNS:localhost, IP:127.0.0.1)的证书,输出到certs/目录。然后创建 Secret:
kubectl create secret generic kafka-tls-certs \ --from-file=ca-cert.pem=./certs/ca-cert.pem \ --from-file=kafka-0-cert.pem=./certs/kafka-0-cert.pem \ --from-file=kafka-0-key.pem=./certs/kafka-0-key.pem \ --from-file=kafka-1-cert.pem=./certs/kafka-1-cert.pem \ --from-file=kafka-1-key.pem=./certs/kafka-1-key.pem \ -n db5.3 TLS 客户端连接验证
在启用 TLS 后,客户端需要携带 CA 证书建立信任:
# 导出 CA 证书 kubectl get secret kafka-tls-certs -n db -o jsonpath='{.data.ca-cert\.pem}' | base64 -d > ca-cert.pem # 客户端属性文件 cat > client.properties << EOL security.protocol=SSL ssl.truststore.location=ca-cert.pem ssl.truststore.type=PEM EOL # 在测试 Pod 内验证(bootstrap 指向 SSL Service) kubectl run kafka-client --rm -it --image=confluentinc/cp-kafka:latest --namespace db -- bash kafka-topics --bootstrap-server kafka-ssl.db.svc.cluster.local:9094 \ --command-config client.properties --list六、部署实战:三种典型场景
以下命令均以 Chart 目录(scripts/helmcharts/databases/charts/kafka)为执行上下文,默认部署到db命名空间(可用--create-namespace自动创建)。
6.1 基础部署(无 TLS)
helm install kafka . --namespace db --create-namespace # 状态检查 kubectl get statefulset -n db kubectl get pods -n db -l app.kubernetes.io/name=kafka kubectl get svc -n db -l app.kubernetes.io/name=kafka安装完成后,NOTES.txt会打印访问地址与集群概要(副本数、KRaft 状态、Cluster ID、TLS 状态)。
6.2 TLS 部署
helm install kafka . -f values-tls.yaml --namespace dbvalues-tls.yaml 是在默认值基础上只做最小改动:启用tls.enabled: true、listeners.ssl.enabled: true、tls.clientAuth: "required",其余保持默认。
6.3 自定义配置
helm install kafka . \ --set replicaCount=3 \ --set persistence.size=200Gi \ --set resources.limits.memory=4Gi \ --set resources.requests.memory=2Gi \ --namespace db6.4 渲染验证与卸载
# 仅渲染模板,不实际部署 helm template test-kafka . helm template test-kafka . -f values-tls.yaml # 卸载(注意:不会删除 PVC) helm uninstall kafka --namespace db # 如需连数据一起清理 kubectl delete pvc -n db -l app.kubernetes.io/name=kafka实现摘要中还记录了从 databases 父 Chart 验证整体渲染的命令cd ../../../ && make db-template(从 chart 目录回到scripts/helmcharts执行);当前仓库的 scripts/helmcharts/Makefile 提供有template(渲染 openreplay 全量模板)、install、clean等目标可供参考。
6.5 扩缩容与监控排查
# 扩缩容(升级方式) helm upgrade kafka . --namespace db --set replicaCount=3 # 监控 kubectl get pods -n db -w -l app.kubernetes.io/name=kafka kubectl logs -n db kafka-0 -f kubectl describe pod kafka-0 -n db kubectl get pvc -n db -l app.kubernetes.io/name=kafka # 连通性排查 kubectl run test-pod --rm -it --image=busybox --namespace db -- sh # 容器内: nc -zv kafka.db.svc.cluster.local 9092 nc -zv kafka-0.kafka-headless.db.svc.cluster.local 9092七、与 databases 父 Chart 的集成
该 Kafka Chart 已作为依赖挂载到 databases 聚合 Chart 中。在 scripts/helmcharts/databases/Chart.yaml 中:
dependencies: - name: kafka repository: file://charts/kafka version: 11.8.6 condition: kafka.enabled父 Chart 的 values.yaml 中维护着同构的kafka:配置块(默认enabled: false),并在 Chart 渲染时透传给子 Chart。这意味着:
- 在
scripts/helmcharts/databases/values.yaml中将kafka.enabled置为true,并按需调整kafka.replicaCount、kafka.persistence.size、kafka.tls等; - Kafka 配置段与父 Chart 的
kafka:段保持字段一一对应; - 通过
make db-template或helm template验证整体渲染。
此外,上游 OpenReplay 应用 Chart(scripts/helmcharts/openreplay/templates/job.yaml 第 551–595 行)在升级时会启动一个 Kafka 迁移 Job(使用kafkaMigration.image,即ghcr.io/openreplay/kafka:3),通过KAFKA_HOST、KAFKA_PORT、KAFKA_SSL、REPLICATION_FACTOR、RETENTION_TIME等环境变量连接 Kafka 并执行 topic 创建/保留策略等操作;这些变量的默认值可在 scripts/helmcharts/vars.yaml 中看到:kafkaHost: "kafka.db.svc.cluster.local"、kafkaPort: "9092"、kafkaUseSsl: "false"、replicationFactor: "2"——即默认期望集群部署在db命名空间,客户端通过 Client Service 域名接入。
八、实现验证与下一步
实现摘要中记录的已验证能力包括:StatefulSet 创建、三类 Service 创建、环境变量渲染、TLS init container、卷挂载、资源限制、控制器仲裁投票生成——均可通过helm template输出清单逐一核对。
如果你要在此基础上继续演进,可以从以下几个方面着手:
- 生产化加固:将
replicaCount提升至 3+ 并相应调高defaultReplicationFactor;为persistence.storageClass指定高性能存储类;启用 TLS 并为各组件下发 CA 信任; - 监控告警:接入 Prometheus/JMX exporter,围绕 QUICKSTART.md 的生产清单建立告警与备份策略;
- 深度定制:借助
extraEnvVars注入任意KAFKA_CFG_*变量,借助extraVolumes/extraVolumeMounts挂载自定义配置或监控探针。
通过本文的源码级拆解可以看到,这个 Chart 的价值不在于堆砌功能,而在于将两份原生清单中的全部关键行为——KRaft 仲裁、Pod 级通告地址、TLS 证书按 Pod 分发——以 Helm 模板函数和参数化 values 的方式完整保留,并做到了随副本数自动伸缩,为 OpenReplay 自托管部署提供了开箱即用、可安全加密的消息基础设施。
【免费下载链接】openreplaySession 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),仅供参考