基于 Terraform 与 Bitnami Kafka Helm Chart 在 GKE 上为 Apache Beam 测试基础设施部署 Kafka 集群
【免费下载链接】beamApache 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 EngineGKE集群上以纯声明式方式落地一套带外部访问能力的 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 Operator01-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-enginekubectl 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.kafkaKafka 集群本体与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.kafkaKafka 集群的完整配置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 分配独立的 LoadBalancerexternalAccess.autoDiscovery.enabledtrue启用自动发现Chart 自动探测节点与端口rbac.createtrue由 Chart 自动创建所需 RBAC 资源service.annotations{networking.gke.io/load-balancer-type: Internal}内部 Service 使用 GKE 内网负载均衡器externalAccess.service.broker.ports.external9094外部访问 Broker 端口为 9094externalAccess.service.controller.containerPorts.external9094Controller 外部容器端口为 9094externalAccess.controller.service.loadBalancerAnnotations3 条 Internal 注解Controller 外部 Service 全部为内网 LBexternalAccess.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的 Deploymentresource 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 falsehelm_release与kubernetes_deployment均不等待 rolloutTerraform apply 会较快返回Pod 的实际就绪由集群侧异步完成。正确做法是等待一段时间后通过 kubectl 观察 Pod 状态直至节点完成扩缩容。调试与排错使用 kafka-client部署完成后可以使用内置的 kafka-client 完成全套验证。查询 kafka-client Pod 名称kubectl get po -l appkafka-client输出类似NAME READY STATUS RESTARTS AGE kafka-client-cdc7c8885-nmcjc 1/1 Running 0 4m12s进入容器 Shellkubectl 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同时验证了客户端到集群的连接链路是否畅通是最快的连通性检查手段。创建 Topickafka-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 集群。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam 测试基础设施指南用 Terraform Helm Provider 部署 Bitnami Kafka 集群Apache Beam 测试基础设施指南用 Terraform Helm Provider 部署 Bitnami Kafka 集群 Apache Beam 仓大数据批处理流处理数据工程Apache Beam 测试基础设施使用 Strimzi 与 Kustomize 在 GKE 上部署持久化 Kafka 集群Apache Beam 测试基础设施使用 Strimzi 与 Kustomize 在 GKE 上部署持久化 Kafka 集群 导读 Apache Beam 的大数据批处理流处理数据工程Apache Beam 测试基础设施 Terraform 指南在 GCP 上自动化部署 GKE 测试集群Apache Beam 测试基础设施 Terraform 指南在 GCP 上自动化部署 GKE 测试集群 本指南以 Apache Beam 仓库中 .test大数据批处理流处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考