news 2026/9/29 20:31:33

基于 Terraform 与 Bitnami Kafka Helm Chart 在 GKE 上为 Apache Beam 测试基础设施部署 Kafka 集群

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于 Terraform 与 Bitnami Kafka Helm Chart 在 GKE 上为 Apache Beam 测试基础设施部署 Kafka 集群

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

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

导读

本文介绍 Apache Beam 仓库中.test-infra/kafka/bitnami模块的完整用法:它借助 Terraform 的 Helm Provider 直接驱动 Bitnami Kafka Helm Chart,在 Google Kubernetes Engine(GKE)集群上以纯声明式方式落地一套带外部访问能力的 Kafka 集群,并为集群内置一个用于验证与排错的 kafka-client Pod。读完本文,你将掌握该模块的前置条件、Terraform 配置的每个关键参数、标准部署流程、GKE Autopilot 环境下的注意事项,以及通过 kafka-client 容器执行kafka-topics.sh、kafka-cluster.sh等命令的完整调试方法。

模块定位:测试基础设施中的 Kafka 部署方案之一

在 Apache Beam 的测试基础设施中,.test-infra/kafka目录集中管理着为集成测试提供 Kafka 环境的各种实现方式,其顶层 README 明确指出该目录下的子目录分别聚焦不同的 Kafka 实现。除本文介绍的 Bitnami 方案外,仓库中还存在另外两条路线:

  • Strimzi 方案:通过 Strimzi Operator(01-strimzi-operator)配合Kafka自定义资源(02-kafka-persistent)来管理 Kafka 集群;
  • 手工 Kubernetes Manifest 方案:基于 Yolean kubernetes-kafka 项目思路,直接以 YAML 编排 3 个 Kafka 副本与 3 个 Zookeeper 副本,通过setup-cluster.sh部署。

而.test-infra/kafka/bitnami模块选择的是"Terraform + Helm Chart"的路线:你不需要在机器上安装 helm 客户端,Terraform 通过其 helm provider 直接在 API 层面完成 Chart 的 release 创建,这使整个 Kafka 环境的生命周期可以被纳入统一的 Terraform 状态管理,便于与 Beam 的 CI/测试流水线集成。

前置要求

在应用该模块前,需要准备以下环境:

要求说明
Terraform安装 Terraform CLI(仓库 .test-infra/kafka/README.md 要求 v1.2.0 及以上)
Kubernetes 集群连接可用的 kubeconfig,能连接到一个 Kubernetes 集群;仓库提供了对应的 GKE 集群 Terraform 模块:.test-infra/terraform/google-cloud-platform/google-kubernetes-engine
kubectl CLI用于后续的调试与排错操作

其中 Kubernetes 集群的获取方式在本仓库内有现成参考:GKE 模块会部署一个私有 GKE 集群,在apache-beam-testing项目下可直接使用us-central1.apache-beam-testing.tfvars/us-west1.apache-beam-testing.tfvars变量文件执行terraform init && terraform apply(详见 google-kubernetes-engine/README.md)。

工作原理:Terraform Helm Provider 直驱 Bitnami Chart

该模块只包含两个 Terraform 文件,逻辑非常聚焦:

  • provider.tf:声明kubernetes与helm两个 Provider,二者都通过~/.kube/config读取集群连接信息;
  • kafka.tf:定义两个核心资源——helm_release.kafka(Kafka 集群本体)与kubernetes_deployment.kafka_client(调试客户端)。

provider.tf的关键内容如下:

provider "kubernetes" { config_path = "~/.kube/config" } provider "helm" { kubernetes { config_path = "~/.kube/config" } }

helm provider 内部直接通过 Kubernetes API 与 Tiller(或 Helm v3 的 release 存储机制)交互,因此在本地机器上无需安装 helm 二进制,也无需单独执行helm install。

helm_release.kafka:Kafka 集群的完整配置

kafka.tf 中的helm_release.kafka是模块的核心,它从https://charts.bitnami.com/bitnami仓库拉取kafkaChart,以kafka为 release 名发布。其中wait = false表示 Terraform 不等待 Pod 全部就绪即返回,这与下文 GKE Autopilot 的"Unschedulable"现象有直接关系。

Chart 的核心参数通过set/set_list注入,汇总如下:

参数值作用
listeners.client.protocolPLAINTEXT客户端监听协议设为明文,无认证加密
listeners.interbroker.protocolPLAINTEXTBroker 间通信协议设为明文
listeners.external.protocolPLAINTEXT外部监听协议设为明文
externalAccess.enabledtrue开启外部访问,为每个 Broker 分配独立的 LoadBalancer
externalAccess.autoDiscovery.enabledtrue启用自动发现(Chart 自动探测节点与端口)
rbac.createtrue由 Chart 自动创建所需 RBAC 资源
service.annotations{"networking.gke.io/load-balancer-type": "Internal"}内部 Service 使用 GKE 内网负载均衡器
externalAccess.service.broker.ports.external9094外部访问 Broker 端口为 9094
externalAccess.service.controller.containerPorts.external9094Controller 外部容器端口为 9094
externalAccess.controller.service.loadBalancerAnnotations3 条 Internal 注解Controller 外部 Service 全部为内网 LB
externalAccess.broker.service.loadBalancerAnnotations3 条 Internal 注解Broker 外部 Service 全部为内网 LB

从配置可以看出该模块面向的是测试场景:全部监听器使用PLAINTEXT(不启用 TLS 与 SASL,降低测试环境的复杂度),同时通过networking.gke.io/load-balancer-type: Internal注解将外部访问限制在 GKE 集群所在的 VPC 内网,避免把 Kafka 暴露到公网。

值得注意的是set_list中externalAccess.controller.service.loadBalancerAnnotations与externalAccess.broker.service.loadBalancerAnnotations各包含 3 条完全相同的注解,这与 Bitnami Kafka Chart 默认按 3 个副本(Broker/Controller 各 3 个)生成外部 Service 的默认行为一一对应——每条注解对应一个副本的 LoadBalancer。

kubernetes_deployment.kafka_client:内置调试客户端

模块同时部署了一个名为kafka-client的 Deployment:

resource "kubernetes_deployment" "kafka_client" { wait_for_rollout = false metadata { name = "kafka-client" labels = { app = "kafka-client" } } spec { selector { match_labels = { app = "kafka-client" } } template { metadata { labels = { app = "kafka-client" } } spec { container { name = "kafka-client" image = "bitnami/kafka:latest" image_pull_policy = "IfNotPresent" command = ["/bin/bash"] args = [ "-c", "while true; do sleep 2; done", ] } } } } }

该 Pod 使用最新的bitnami/kafka:latest镜像,容器启动后进入一个while true; do sleep 2; done的空循环保持存活,专门用于在集群内验证连接、创建 Topic、查询元数据等排错操作——这意味着部署完 Kafka 后,你无需再额外准备任何客户端机器即可开展调试。

标准部署流程

该模块遵循标准 Terraform 工作流,在.test-infra/kafka/bitnami目录下依次执行:

terraform init terraform apply

前提是你的 kubeconfig(~/.kube/config)已经指向目标集群。若需要结合仓库中的 GKE 模块,可先按 google-kubernetes-engine/README.md 的流程创建集群并配置好 kubeconfig,再回到本模块执行上述两条命令。

GKE Autopilot 下的特殊注意事项

当把该模块应用到 GKE Autopilot 集群时,你会观察到 Kafka 相关 Pod 长期处于"Unschedulable"(不可调度)状态。README 明确解释了原因:Autopilot 集群需要时间扩容节点,Kubernetes 只有在计算资源就绪后才会真正调度并完成 Kafka 集群的创建。

因此遇到该状态不必恐慌,也不应立即判定部署失败——结合 kafka.tf 中两处wait = false(helm_release与kubernetes_deployment均不等待 rollout),Terraform apply 会较快返回,Pod 的实际就绪由集群侧异步完成。正确做法是等待一段时间后通过 kubectl 观察 Pod 状态,直至节点完成扩缩容。

调试与排错:使用 kafka-client

部署完成后,可以使用内置的 kafka-client 完成全套验证。

查询 kafka-client Pod 名称

kubectl get po -l app=kafka-client

输出类似:

NAME READY STATUS RESTARTS AGE kafka-client-cdc7c8885-nmcjc 1/1 Running 0 4m12s

进入容器 Shell

kubectl exec --stdin --tty kafka-client-cdc7c8885-nmcjc -- /bin/bash

容器基于bitnami/kafka:latest镜像构建,所有必需的kafka-*.sh脚本已在其 PATH 中,无需额外安装任何客户端工具。

关键连接参数:--bootstrap-server kafka:9092

所有 Kafka 命令都可以直接使用:

--bootstrap-server kafka:9092

原因是 kafka-client Pod 与 Kafka 集群位于同一 Kubernetes 集群中,可利用集群内置 DNS 服务解析到名为kafka的 Service——这正是 Bitnami Helm 操作符创建的、暴露 9092 端口的 Kubernetes Service(对应 kafka.tf 中listeners.client.protocol = PLAINTEXT的客户端监听器)。

获取集群 ID 并验证连接

kafka-cluster.sh cluster-id --bootstrap-server kafka:9092

该命令返回集群 ID,同时验证了客户端到集群的连接链路是否畅通,是最快的连通性检查手段。

创建 Topic

kafka-topics.sh --create --topic some-topic --partitions 3 --replication-factor 3 --bootstrap-server kafka:9092

此处--partitions 3与--replication-factor 3与该模块按 3 副本部署的集群拓扑相匹配——复制因子 3 意味着每个分区在 3 个 Broker 上各存一份副本。

查询 Topic 信息

kafka-topics.sh --describe --topic some-topic --bootstrap-server kafka:9092

该命令输出some-topic的分区数、副本分布、ISR(同步副本)等元数据,可用于验证分区与副本是否按预期分配。

如果希望进一步了解在容器内执行命令的通用方法论,可参考 Kubernetes 官方文档 中关于进入运行中容器的说明。

结合 Beam 使用场景的延伸

这套 Bitnami Kafka 集群在仓库中服务于 Beam 的 Kafka IO 相关测试与示例环境,例如sdks/java/io/kafka的集成测试、it/kafka目录下的测试用例,以及learning/tour-of-beam/io中的 Kafka 学习内容,都可能需要这样一个集群作为后端。与 Strimzi 方案 和 手工 Manifest 方案 相比,本模块的优势在于:以 Terraform 单一工具链完成"集群创建 + Chart 发布 + 调试客户端部署"的全部环节,无需在本地维护 helm 客户端,且内网负载均衡的配置天然适配 GKE 测试环境的网络安全边界。

总结

.test-infra/kafka/bitnami模块为 Apache Beam 测试基础设施提供了一条低门槛、可复现的 Kafka 环境交付路径:用两个 Terraform 文件完成从 Chart 发布到调试客户端的全链路部署,内置的kafka-client让验证与排错不再依赖额外工具。无论是标准 GKE 集群还是 Autopilot,理解其配置参数(监听协议、外部访问、内网 LB)与运行机制(wait = false、DNS 直连kafka:9092)之后,你都能快速搭建并验证一套可用的测试级 Kafka 集群。

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

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

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/29 20:30:33

报错型SQL注入深度解析:updatexml、extractvalue与实战防御

报错型SQL注入是我在漏洞挖掘和靶场通关里最先熟练掌握的一类技巧。相比联合查询需要清点字段数、盲注需要逐位去猜,报错注入只要让数据库把错误信息打到页面上,你就能直接从报错里读到库名、表名、字段名甚至是最终数据。可以这么说:在目标应…

作者头像 李华
网站建设 2026/9/29 20:28:49

​政企采购数字孪生服务商,武汉启创动力 4 个维度选型指南

数字孪生这几年在政企圈很热,但真正落地之后能持续用起来的项目,比例并不高。不少单位花了几十万甚至上百万搭了一套三维可视化系统,验收时大屏效果惊艳,半年后却沦为偶尔接待参观时“点亮一下”的工具。问题出在哪儿?…

作者头像 李华
网站建设 2026/9/29 20:27:43

零基础教程:用 TaoToken 统一 Key 把 AI 模型接入手机端零代码应用

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华