- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
Apache Beam 的持续集成(CI)与集成测试(IT)需要真实可用的消息中间件环境来验证 Kafka I/O 连接器。本文基于仓库中的 .test-infra/kafka/README.md 及其配套子模块,系统讲解如何借助 Kubernetes manifests 与 Terraform 在 Google Kubernetes Engine (GKE) 上按需供给 Kafka 集群。读完本文,你将掌握 Strimzi Operator、Bitnami Helm Chart 与 GCP 私有 IP 代理三种方案的部署流程、关键参数与排障手段,并能直接复用到 Beam 的测试环境搭建中。
目录总览:一个目录,多种 Kafka 实现
.test-infra/kafka/目录的核心设计思想是:按实现方案划分子目录,每个子目录自包含部署所需的 Terraform 或 Kubernetes 清单。这样既可以在不同测试场景间切换后端实现,又不互相干扰。
.test-infra/kafka/ ├── README.md # 总览文档(本文主体) ├── bitnami/ # 基于 Bitnami Kafka Helm Chart 的部署方案 ├── proxy/ # GCP 私有 IP 代理(bastion host)方案 └── strimzi/ # 基于 Strimzi Operator 的部署方案 ├── 01-strimzi-operator/ # 通过 Terraform 部署 Strimzi Operator └── 02-kafka-persistent/ # 通过 Kustomize 部署持久化 Kafka 集群从仓库结构看,.test-infra/下还包含 dataproc、pubsub、jupyter、metrics 等同级测试基础设施模块,Kafka 目录只是其中一环,专门服务于 Kafka 相关集成测试的中间件供给。
环境要求:部署前需要准备什么
根据 .test-infra/kafka/README.md 中的 Requirements 小节,使用这套代码需要满足以下前置条件:
| 依赖项 | 说明 |
|---|---|
| Terraform CLI v1.2.0 及以上 | 所有 Terraform 模块均基于该版本语法编写 |
| GKE 集群 | Ingress 配置默认假设运行在 Google Kubernetes Engine 集群上,集群的创建可参考仓库中的 google-kubernetes-engine 模块 |
| kubectl CLI | 用于执行 Kustomize apply 与集群状态查询 |
| IntelliJ 或 VS Code 的 Kubernetes 插件 | 可选但强烈推荐,便于在本地直接查看和编辑 YAML manifests |
其中 GKE 集群模块位于 .test-infra/terraform/google-cloud-platform/google-kubernetes-engine,是先决依赖:所有 Kafka 方案都要求先有一个可连接的 Kubernetes 集群(后续的 Strimzi、Bitnami 子文档也重复强调了这一点)。
通用使用方式:聚焦单一实现的模块化布局
总览文档指出,每个子目录聚焦一种特定的 Kafka 实现,并给出了以 strimzi 目录为例的用法。这意味着选择方案 = 进入对应子目录,按该子目录自带的 README 执行,而不是在根目录统一执行。后文将逐一展开三种方案的完整步骤。
方案一:Strimzi —— Operator + 持久化集群两段式部署
Strimzi 是 Kubernetes 上最主流的 Kafka Operator 之一。Beam 测试仓库将其拆为两个阶段:先用 Terraform 部署 Operator,再用 Kustomize 部署由 Operator 管理的持久化 Kafka 集群。两阶段均有独立文档与完整配置。
阶段一:用 Terraform Helm Provider 部署 Strimzi Operator
.test-infra/kafka/strimzi/README.md 说明该模块通过 Helm Chart 形式安装 Strimzi Operator,其核心实现位于 kafka.tf,使用helm_release资源完成部署:
resource "helm_release" "strimzi-helm-release" { name = var.name namespace = var.namespace create_namespace = true repository = var.chart_repository chart = var.chart_name version = var.chart_version wait = false set { name = "watchAnyNamespace" value = "true" } }几个值得注意的实现细节:
wait = false:Terraform 不会等待 Helm release 完全就绪才返回,加快了 apply 速度,后续可通过 kubectl 确认状态;watchAnyNamespace = true:Operator 会监听整个集群的所有命名空间,从而可以管理后续在strimzi命名空间之外创建的 Kafka 自定义资源;create_namespace = true:目标命名空间不存在时自动创建。
模块暴露的全部变量定义在 variables.tf 中,而默认值集中在 common.tfvars:
name = "strimzi" namespace = "strimzi" kubeconfig_path = "~/.kube/config" chart_name = "strimzi-kafka-operator" chart_version = "0.40.0" chart_repository = "https://strimzi.io/charts/"各变量的语义与默认值如下:
| 变量 | 默认值 | 含义 |
|---|---|---|
name | strimzi | Helm release 名称 |
namespace | strimzi | Operator 部署的命名空间 |
kubeconfig_path | ~/.kube/config | 连接集群用的 kubeconfig 路径,由 provider.tf 中的 Helm Provider 读取 |
chart_name | strimzi-kafka-operator | 使用的 Helm Chart 名称 |
chart_version | 0.40.0 | Chart 版本(按当前仓库配置为准) |
chart_repository | Strimzi 官方 Helm 仓库 | Chart 来源仓库地址 |
执行标准 Terraform 工作流即可完成 Operator 部署(沿用 common.tfvars 的默认配置):
DIR=.test-infra/kafka/strimzi/01-strimzi-operator VARS=common.tfvars # 注意该文件位于 $DIR 下,使用相对文件名即可 terraform -chdir=$DIR init terraform -chdir=$DIR apply -var-file=$VARS-chdir让 Terraform 在指定目录内执行,而-var-file传入变量文件路径;这里刻意使用相对文件名,正是为了让两个命令在任何工作目录下都能拼出正确路径。
阶段二:用 Kustomize 部署持久化 Kafka 集群
Operator 就绪后,真正意义上的 Kafka 集群由 02-kafka-persistent 模块提供,其文档 README.md 明确了集群清单源自 Strimzi 官方kafka-persistent.yaml示例(v0.33.2 版本)的再分发,目录结构采用标准的base + overlays组织:
02-kafka-persistent/ ├── base/ │ └── v0.33.2/ │ ├── kafka-persistent.yaml # 集群主清单 │ └── kustomization.yaml └── overlays/ └── gke-internal-load-balanced/ ├── kustomization.yaml # 引用 base 并叠加补丁 └── listeners.yaml # 补丁:GKE 内部 TCP LoadBalancerbase 层定义集群本体。kafka-persistent.yaml 的核心配置如下:
apiVersion: kafka.strimzi.io/v1beta2 kind: Kafka metadata: name: beam-testing-cluster spec: kafka: resources: requests: cpu: 8 memory: 64Gi version: 3.6.0 replicas: 3 config: offsets.topic.replication.factor: 3 transaction.state.log.replication.factor: 3 transaction.state.log.min.isr: 2 default.replication.factor: 3 min.insync.replicas: 2 inter.broker.protocol.version: "3.4" storage: type: jbod volumes: - id: 0 type: persistent-claim size: 500Gi deleteClaim: false zookeeper: resources: requests: cpu: 1 memory: 2Gi replicas: 3 storage: type: persistent-claim size: 100Gi deleteClaim: false entityOperator: topicOperator: {} userOperator: {}对照配置逐项解读:
- 集群规模:
replicas: 3的 Kafka broker 与 3 副本 Zookeeper,配合entityOperator(topicOperator + userOperator)组成完整管控面; - Kafka 版本:
version: 3.6.0,inter.broker.protocol.version: "3.4"保证 broker 间协议兼容性,这是升级场景中的关键参数; - 可靠性配置:
offsets.topic.replication.factor、transaction.state.log.replication.factor、default.replication.factor均为 3,min.insync.replicas与transaction.state.log.min.isr为 2 —— 保证任意单 broker 故障下集群仍可用,同时满足事务与 Exactly-Once 语义的 ISR 要求; - 存储:Kafka 使用
jbod类型挂载 500Gi 的 PersistentVolumeClaim,Zookeeper 使用 100Gi 的持久化卷,且deleteClaim: false意味着删除集群 CR 时 PVC 会被保留(避免误删数据); - 资源申请:Kafka 每 broker 申请 8 核 / 64Gi 内存、Zookeeper 每节点 1 核 / 2Gi,这是为 Beam 大规模集成测试准备的余量,生产化时可按需下调。
overlays 层则负责适配目标环境。从 overlays/gke-internal-load-balanced/kustomization.yaml 可以看出,该 overlay 引用../../base/v0.33.2作为基础,并通过listeners.yaml补丁为集群追加一个GKE 内部 TCP LoadBalancer 类型的 Ingress,使集群只能被 VPC 内的测试执行者访问,而不会暴露到公网——这正是测试环境的安全基线。
部署命令(使用 kubectl 内建的 Kustomize 插件):
kubectl apply -k .test-infra/kafka/strimzi/02-kafka-persistent/overlays/gke-internal-load-balanced --namespace=strimzi等待集群进入 Ready 状态(超时 1200 秒,因为需要为 PVC 和节点容量等待调度):
kubectl wait kafka beam-testing-cluster --for=condition=Ready --timeout=1200s获取连接地址:逐副本查询 LoadBalancer IP 与端口
Strimzi 为每个 broker 副本各生成一个独立的 Service,命名规则为beam-testing-cluster-kafka-<副本序号>,默认副本序号为 1 到 3。部署文档给出了用jsonpath精确抽取连接信息的命令:
# 获取第 N 个副本的负载均衡器 IP kubectl get svc beam-testing-cluster-kafka-$REPLICA_NUMBER -o jsonpath='{.status.loadBalancer.ingress[0].ip}' # 获取第 N 个副本对外暴露的端口 kubectl get svc beam-testing-cluster-kafka-$REPLICA_NUMBER -o jsonpath='{.spec.ports[0].port}'其中$REPLICA_NUMBER取 1、2、3。将这些 IP:Port 组合填入 Beam Kafka I/O 的bootstrap.servers,即可让测试管道连入集群。
方案二:Bitnami —— 基于 Helm Chart 与 Terraform 的轻量方案
如果希望避开 Kustomize/Operator 的复杂概念,仓库提供了基于 Bitnami Kafka Helm Chart 的第二方案,见 bitnami/README.md。
该方案的最大特点在于:模块内部直接使用 Terraform 的 Helm Provider,因此你本机甚至不需要安装 helm CLI。部署遵循最标准的 Terraform 工作流:
terraform init terraform applyGKE Autopilot 下的特殊注意点
文档特别提示:当部署目标是 GKE Autopilot 集群时,Pod 会先呈现Unschedulable状态。这不是故障,而是Autopilot 需要时间自动扩容节点;等节点资源就绪后,Kubernetes 会自动完成 Kafka 集群的调度与启动。遇到该状态时耐心等待即可。
内置 Kafka 客户端容器用于调试
该模块会在集群内额外部署一个kafka-clientPod(基于最新 Bitnami Kafka 镜像),内置全部kafka-*.sh脚本,专供排障使用。
1. 查询客户端 Pod 名称:
kubectl get po -l app=kafka-client预期输出类似:
NAME READY STATUS RESTARTS AGE kafka-client-cdc7c8885-nmcjc 1/1 Running 0 4m12s2. 进入容器 shell:
kubectl exec --stdin --tty kafka-client-cdc7c8885-nmcjc -- /bin/bash3. 执行 Kafka 命令:由于客户端与集群同处一个 Kubernetes 集群,直接借助集群 DNS 使用--bootstrap-server kafka:9092(Bitnami 部署会创建一个名为kafka的 Service,暴露 9092 端口):
# 获取 cluster-id,同时验证连通性 kafka-cluster.sh cluster-id --bootstrap-server kafka:9092 # 创建 topic(3 分区、3 副本) kafka-topics.sh --create --topic some-topic --partitions 3 --replication-factor 3 --bootstrap-server kafka:9092 # 查看 topic 详情 kafka-topics.sh --describe --topic some-topic --bootstrap-server kafka:9092这套"集群内客户端"模式非常适合 CI 场景:测试管道可以在集群内直接完成 topic 预创建、消息注入与消费验证,无需从外部网络访问 broker。
方案三:proxy —— 为私有 IP Kafka 打通外部访问通道
Strimzi 的 GKE 内部负载均衡器方案决定了集群只对 VPC 内可见。当 Beam 测试执行者位于集群外部(如本地开发机)时,仓库提供了 proxy 模块:在 Google Cloud 上创建一个私有 IP 的 bastion host(跳板机),作为访问私有 Kafka 实例的代理,详见 proxy/README.md。
前置条件
使用该模块的前提是已有一套 Kafka 集群(即前文 Strimzi 或 Bitnami 方案的产出,参见 .test-infra/kafka)。
第一步:准备 bootstrap server 映射
模块的关键变量是bootstrap_endpoint_mapping,用于把Kafka bootstrap server 主机名映射到希望代理暴露的本地端口,具体定义见 variables.tf。即告诉代理:"<broker-ip>:9092这个地址,请在本地用 9092 端口转发"。
第二步:应用模块
该模块不使用 Terraform backend(状态文件保存在本地),按典型工作流执行:
DIR=.test-infra/kafka/proxy terraform -chdir=$DIR initterraform -chdir=$DIR apply -var-file=common.tfvars -var-file=name_of_your_specific.tfvars注意这里需要同时传入两个变量文件:common.tfvars提供通用默认值(位于 .test-infra/kafka/proxy/common.tfvars),后一个是你为本次集群准备的专属 tfvars。
第三步:使用 gcloud SSH 隧道
模块 apply 成功后会直接输出一条可复制的 gcloud 隧道命令,格式形如:
gcloud compute ssh yourinstance --tunnel-through-iap --project=project --zone=zone --ssh-flag="-4 -L9093:localhost:9093" --ssh-flag="-4 -L9092:localhost:9092" --ssh-flag="-4 -L9094:localhost:9094"--tunnel-through-iap通过 Identity-Aware Proxy 建立安全隧道,-L参数把本地端口逐一映射到代理上的 Kafka 端口。执行后,本地localhost:9092即可作为bootstrap.servers供 Beam 管道使用,且全程不经公网直接暴露 Kafka。
三种方案如何选型
结合三份子文档,可以给出如下选型建议(基于仓库现状的合理推断):
| 对比维度 | Strimzi(推荐主路径) | Bitnami | proxy(配套组件) |
|---|---|---|---|
| 部署方式 | Terraform + Helm(Operator),Kustomize(集群) | 纯 Terraform Helm Provider | Terraform 计算实例 |
| 集群能力 | 持久化、3 副本、事务参数完整 | 开箱即用的 Bitnami 默认配置 | 不建集群,仅打通访问 |
| 适用场景 | Beam 大规模、长时运行、需持久化的集成测试 | 快速起停、轻量验证、Autopilot 环境 | 外部执行者访问私有集群时使用 |
| 调试手段 | kubectl 查 CR 状态 / Service | 内置 kafka-client 容器 | gcloud SSH 隧道 |
从仓库结构看,Strimzi 方案被组织得最完整(独立 operator 模块 + base/overlays 两层 Kustomize),是 Beam 测试基础设施的主推实现;Bitnami 方案作为轻量备选;proxy 模块则是为跨网络访问场景补齐的最后一块拼图。三者共同构成了从"集群供给"到"外部接入"的完整链路。
结语
本文围绕 .test-infra/kafka/README.md 的骨架,结合仓库内各子模块的 Terraform 配置与 Kubernetes 清单,完整还原了 Beam 测试基础设施中 Kafka 集群的供给全流程:Strimzi 的 Operator + 持久化集群两段式部署(含 3 副本、事务参数、500Gi 存储等生产级配置细节)、Bitnami 的轻量 Helm 方案(含 Autopilot 注意事项与内置 kafka-client 调试容器)、以及 proxy 的 IAP 隧道外部接入方案。读者可依据测试规模与网络环境直接套用对应命令,在 GKE 上快速复现这套 Kafka 测试环境。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 测试基础设施中的 Strimzi Kafka Operator 部署指南:Terraform + Helm + Kustomize 全流程
Apache Beam 测试基础设施中的 Strimzi Kafka Operator 部署指南:Terraform + Helm + Kustomize 全流
大数据批处理流处理数据工程Facebook-Messenger-Bot入门:30分钟快速搭建你的个性化AI聊天助手
Facebook Messenger Bot入门:30分钟快速搭建你的个性化AI聊天助手 Facebook Messenger Bot是一个基于Seq2Seq模
大数据批处理流处理数据工程Apache Beam 测试基础设施的 GCP Terraform IaC 实战指南:从私有 GKE 集群到 Vertex AI Featurestore
Apache Beam 测试基础设施的 GCP Terraform IaC 实战指南:从私有 GKE 集群到 Vertex AI Featurestore Ap
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考