- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本篇指南聚焦 Apache Beam 仓库.test-infra/kafka/bitnami模块:它如何借助 Bitnami Kafka Helm Chart 与 Terraform Helm Provider,在不安装 helm 客户端的前提下,将一套多 Broker Kafka 集群一键部署到(Google)Kubernetes 集群,并同时拉起一个专用的kafka-client调试容器。读完本文,你将掌握该模块的完整架构、部署流程、GKE Autopilot 下的注意事项,以及通过 kubectl 进入客户端容器执行kafka-cluster.sh、kafka-topics.sh等命令进行连通性验证与 Topic 管理的完整排障方法。
模块定位:Beam 测试基础设施中的 Kafka 供应方案
Apache Beam 的持续集成与测试体系需要在真实环境下验证各类 IO 连接器,Kafka 就是其中之一——it/kafka下的集成测试(例如 KafkaIOLT.java、KafkaIOST.java)依赖可用的 Broker 端点,KafkaResourceManager.java 甚至直接以PLAINTEXT://host:port形式拼装 bootstrap 连接串(见其KAFKA_BROKER_PORT相关逻辑)。
仓库的 .test-infra/kafka/README.md 说明:该目录集中存放用于供应 Kafka 集群的 Kubernetes/Terraform 清单,并按实现方式划分为多个子模块:
bitnami:基于 Bitnami 官方 Helm Chart 部署(本文主角);strimzi:基于 Strimzi Kafka Operator 部署(01-strimzi-operator部署 Operator,02-kafka-persistent提供持久化 Kafka 集群清单);proxy:在 GCP 上供应一台私有 IP 堡垒机,作为私有 Kafka 实例的代理入口。
.test-infra/kafka/bitnami是其中最"开箱即用"的一种:它直接消费 Bitnami 维护的 kafka Helm Chart,把集群的编排细节全部交给 Chart,模块本身只负责通过 Terraform 声明所需的自定义配置。
工作原理:为什么不需要安装 helm
该模块的独特之处在于:它使用的是 Terraform Helm Provider(官方文档见 registry.terraform.io/providers/hashicorp/helm),因此你无需在本地或 CI 环境中安装 helm CLI 即可完成 Chart 的渲染与下发。Helm Chart 的解析、模板渲染与 release 管理全部由 Terraform 的 helm provider 在后台完成。
这一点可以从 provider.tf 得到印证:模块同时声明了kubernetes与helm两个 provider,且都读取本地的~/.kube/config作为集群连接凭据:
provider "kubernetes" { config_path = "~/.kube/config" } provider "helm" { kubernetes { config_path = "~/.kube/config" } }也就是说,只要你的kubectl能访问目标集群,Terraform 就能直接向该集群部署 Helm Release——这也是整个模块所有部署操作的前置基础。
模块清单剖析:从 kafka.tf 看集群配置
模块的核心配置位于 kafka.tf,其中包含两个资源:helm_release.kafka(Kafka 集群本体)和kubernetes_deployment.kafka_client(调试客户端)。
helm_release.kafka:集群本体
resource "helm_release" "kafka" { wait = false repository = "https://charts.bitnami.com/bitnami" chart = "kafka" name = "kafka" ... }关键点与取值说明:
- Chart 来源:
repository = "https://charts.bitnami.com/bitnami"、chart = "kafka",即使用 Bitnami 官方 Chart 仓库中的kafkaChart,release 名为kafka; wait = false:Terraform 在 release 创建后不等待所有 Pod 就绪,避免因节点扩容等异步过程阻塞 apply(这一点与后文 GKE Autopilot 的"Unschedulable"现象直接相关);- 监听协议统一为 PLAINTEXT:通过
set将listeners.client.protocol、listeners.interbroker.protocol、listeners.external.protocol全部设为PLAINTEXT。测试环境不启用 TLS/SASL,简化了客户端连接(KafkaResourceManager.java 中拼接的PLAINTEXT://连接串与此一致); - 外部访问:
externalAccess.enabled = true、externalAccess.autoDiscovery.enabled = true,由 Chart 自动为每个 Broker/Controller 创建独立的 LoadBalancer Service,外部端口固定为9094(externalAccess.service.broker.ports.external与externalAccess.service.controller.containerPorts.external均为9094); - RBAC:
rbac.create = true,让 Chart 自建所需的 ServiceAccount 与权限; - GKE 内部负载均衡:模块通过
service.annotations与externalAccess.*.service.loadBalancerAnnotations设置networking.gke.io/load-balancer-type: Internal注解(set_list中为 3 个副本各配置一条),确保对外暴露的负载均衡器全部为 GKE 内部 LB,不向公网开放; - 客户端监听端口:Chart 默认创建的
kafkaService 在集群内暴露9092,供 Pod 间与调试容器使用(见下文)。
kubernetes_deployment.kafka_client:内置调试客户端
resource "kubernetes_deployment" "kafka_client" { wait_for_rollout = false metadata { name = "kafka-client" 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", ] } } }该 Deployment 以bitnami/kafka:latest镜像运行一个常驻bash(死循环保活),Pod 打上app=kafka-client标签。镜像内置了全套kafka-*.sh管理脚本,因此它既是部署后的连通性验证工具,也是日常调试与故障排查的"现场环境"。
前置条件(Requirements)
在应用本模块前,需要准备:
- Terraform CLI(terraform.io):仓库根 .test-infra/kafka/README.md 要求 v1.2.0 及以上;
- 一个可达的 Kubernetes 集群:
kubectl能正常访问,且~/.kube/config已配置好凭据。若使用 GCP,可先通过仓库内的 .test-infra/terraform/google-cloud-platform/google-kubernetes-engine 模块供应私有 GKE 集群——该模块要求预先存在 VPC/子网、最小权限的 Service Account 以及用于 Terraform 后端状态的 GCS Bucket,具体步骤见其 README.md; - kubectl CLI:用于后续查看 Pod、进入容器执行 Kafka 命令。
使用步骤(Usage)
本模块遵循标准 Terraform 工作流,仅需两条命令:
terraform init terraform applyterraform init会拉取 helm/kubernetes provider 并解析 Chart 依赖;terraform apply将执行 kafka.tf 中声明的两个资源。若需要指定集群变量(如使用其他 tfvars),可按仓库惯例通过-chdir与-var-file组合,例如参照 strimzi 模块 的用法:terraform -chdir=$DIR apply -var-file=$VARS。
特别提示:GKE Autopilot 下的调度行为
如果你的目标集群是GKE Autopilot模式,apply 之后你会观察到 Kafka Pod 长时间处于"Unschedulable"状态。这并非故障:Autopilot 集群需要时间自动扩容节点池,只有当底层计算资源真正就绪后,Kubernetes 才会完成 Kafka 集群的调度与拉起。因此请耐心等待,不要误判为部署失败而反复重建资源;配合helm_release.kafka的wait = false设置,整个供应过程是异步、宽松的。
调试与故障排查(Debugging and Troubleshooting)
部署完成后,模块自带的kafka-client容器就是你的"调试工作台"。以下操作全部在集群内进行。
1. 查询 kafka-client Pod 名称
kubectl get po -l app=kafka-client输出示例:
NAME READY STATUS RESTARTS AGE kafka-client-cdc7c8885-nmcjc 1/1 Running 0 4m12skafka-client-cdc7c8885-nmcjc即 Deployment 生成的 Pod(名称中的随机后缀由 ReplicaSet 生成)。
2. 进入容器获取 Shell
kubectl exec --stdin --tty kafka-client-cdc7c8885-nmcjc -- /bin/bash此命令打开容器内交互式 bash(详见 Kubernetes 官方文档)。
3. 执行 Kafka 管理命令
容器基于最新的bitnami/kafka镜像,路径中已预装全部kafka-*.sh管理脚本。由于客户端 Pod 与 Kafka 集群位于同一 Kubernetes 集群,所有命令均可直接使用:
--bootstrap-server kafka:9092这是因为 Bitnami Chart 会创建一个名为kafka的 Kubernetes Service,暴露端口9092,集群内 Pod 通过 DNS 即可解析到该 Service 背后的 Broker。
获取 cluster-id(连通性验证):
kafka-cluster.sh cluster-id --bootstrap-server kafka:9092该命令返回集群 ID,同时验证客户端到 Broker 的连通性是否正常——这是排障的第一步,任何配置错误都会在此暴露。
创建 Topic:
kafka-topics.sh --create --topic some-topic --partitions 3 --replication-factor 3 --bootstrap-server kafka:9092创建一个名为some-topic、3 个分区、副本因子为 3 的 Topic(对应 3 副本的典型测试配置,若实际 Broker 数不足 3,请按需下调--replication-factor)。
查看 Topic 信息:
kafka-topics.sh --describe --topic some-topic --bootstrap-server kafka:9092输出该 Topic 的分区分布、Leader、ISR 等元数据,用于核对副本是否均匀分配、集群是否健康。
延伸:从集群内到集群外
若 Beam 集成测试运行在集群之外(例如本地或 CI),需要将 Broker 暴露到外部。本模块已将外部监听端口统一配置为9094,且所有 LoadBalancer 均为 GKE 内部 LB。此时可参考同目录下的 proxy 模块:它会在 GCP 上创建一台私有 IP 堡垒机作为代理,通过bootstrap_endpoint_mapping变量(见 variables.tf)将kubectl get svc得到的各 LoadBalancer IP(如10.1.2.3:9094)映射到本地端口,apply 成功后会输出形如gcloud compute ssh ... --tunnel-through-iap --ssh-flag="-4 -L9094:localhost:9094"的隧道命令,从而把私有 Kafka 流量安全地转发到开发机。
小结
.test-infra/kafka/bitnami模块为 Apache Beam 测试体系提供了一条最快捷的 Kafka 供应路径:Terraform Helm Provider 免去 helm 客户端依赖,Chart 参数化配置(PLAINTEXT 协议、9094 外部端口、内部负载均衡)与自带的kafka-client调试容器构成了"部署即验证"的闭环。结合 kafka.tf 的源码阅读与 it/kafka 集成测试中的连接串约定,你可以快速复制这套方案,为自己的 Beam Kafka 集成测试搭建可复现的 Kafka 环境。
- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
基于 Terraform 与 Bitnami Kafka Helm Chart 在 GKE 上为 Apache Beam 测试基础设施部署 Kafka 集群
基于 Terraform 与 Bitnami Kafka Helm Chart 在 GKE 上为 Apache Beam 测试基础设施部署 Kafka 集群 导
Apache Beam 测试基础设施指南:用 Terraform Helm Provider 部署 Bitnami Kafka 集群
Apache Beam 测试基础设施指南:用 Terraform Helm Provider 部署 Bitnami Kafka 集群 Apache Beam 仓
大数据批处理流处理数据工程Apache Beam 测试基础设施实战:用 Terraform 与 Bitnami Helm Chart 在 Kubernetes 中部署 Redis 集群(02.redis 模块)
Apache Beam 测试基础设施实战:用 Terraform 与 Bitnami Helm Chart 在 Kubernetes 中部署 Redis 集群(
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考