ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam 测试基础设施:使用 Kustomize 在 Kubernetes 上安装 Strimzi Kafka Operator

Apache Beam 测试基础设施:使用 Kustomize 在 Kubernetes 上安装 Strimzi Kafka Operator 【免费下载链接】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/strimzi 目录下的第一份 kustomization——01-strimzi-operator——展开讲解 Beam 官方测试基础设施如何在 Kubernetes 集群中通过 Kustomize 安装 Strimzi Kafka Operator并以此为基底进一步部署持久化 Kafka 集群供 Beam 的 KafkaIO 集成测试使用。读完本文你将掌握该 kustomization 的目录组织与版本管理方式、Operator 各 manifest 组件的构成与关键环境变量以及从安装、等待就绪、端口转发到用kcat/KafkaIO验证连接的一整套可复现操作流程。一、背景Beam 测试环境为什么需要 StrimziApache Beam 的 Java SDK 提供了成熟的 KafkaIO 读写连接器用于在流式与批式管道中消费和写入 Kafka 主题。为了在持续集成环境中验证 KafkaIO 的真实行为Beam 仓库在 .test-infra/kafka/strimzi 目录下维护了一套完整的测试基础设施其职责是在 Kubernetes 上安装 Strimzi Kafka Operator通过 Operator 声明式地创建带持久化存储的 Kafka 集群为本地验证和 Dataflow 上的 KafkaIO 管道提供可访问的 bootstrap server 地址。该目录中的所有资源都使用 Kustomize 管理每个子目录按照必须被应用的顺序编号命名。其中 01-strimzi-operator 是第一步也是最基础的一步它只负责把 Strimzi Operator 本身装进集群此时集群里还没有任何 Kafka broker。二、目录结构与版本化管理方式01-strimzi-operator目录的完整结构如下.test-infra/kafka/strimzi/ ├── 01-strimzi-operator/ │ ├── README.md # 本文关联的核心文档 │ ├── kustomization.yaml # 顶层 kustomization引入 namespace 版本目录 │ ├── namespace.yaml # 创建 strimzi 命名空间 │ └── v0.33.2/ # 按 Strimzi 版本号命名的 manifest 目录 │ ├── kustomization.yaml # 版本内 kustomization列出全部 operator 资源 │ ├── 010-ServiceAccount-*.yaml │ ├── 020-*.yaml / 021-*.yaml / 022-*.yaml / 023-*.yaml # RBAC │ ├── 030-*.yaml / 031-*.yaml / 033-*.yaml # 委托 RBAC │ ├── 040-Crd-kafka.yaml ... 049-Crd-kafkarebalance.yaml # CRD 定义 │ ├── 050-ConfigMap-strimzi-cluster-operator.yaml # 日志配置 │ └── 060-Deployment-strimzi-cluster-operator.yaml # Operator 本体 ├── 02-kafka-persistent/ # 第二步Kafka 持久化集群base overlays └── README.md # 总入口requirements 与完整 usage根据 01-strimzi-operator/README.md 的说明这份 kustomization 安装的 Strimzi Operator 清单是从上游源码重新分发redistributed的原始出处为 Strimzi 项目仓库的install/cluster-operator目录。包含 Operator manifest 的目录按发布版本命名——当前仓库锁定的是v0.33.2与之配套的 Kafka 集群基线 02-kafka-persistent/base/v0.33.2 同样按版本命名二者保持版本一致。需要了解完整的前置条件与使用步骤时文档明确指向总入口 .test-infra/kafka/strimzi/README.md。顶层 kustomization 干了什么顶层 kustomization.yaml 内容非常简短但定义了整体装配关系# Installs the strimzi operator in your kubernetes cluster. namespace: strimzi resources: - namespace.yaml - v0.33.2namespace: strimzi为整个 kustomization 内的所有资源统一设置目标命名空间resources依次引入namespace.yaml与版本目录v0.33.2。其中 namespace.yaml 创建一个名为strimzi的 NamespaceapiVersion: v1 kind: Namespace metadata: name: strimzi而 v0.33.2/kustomization.yaml 则按编号顺序列出该版本 Operator 的全部 28 个资源文件编号即应用顺序先后包括ServiceAccount、ClusterRole/RoleBinding、ClusterRoleBinding、CRDkafka、kafkaconnect、strimzipodset、kafkatopic、kafkauser、kafkamirrormaker、kafkabridge、kafkaconnector、kafkamirrormaker2、kafkarebalance、ConfigMap 与 Deployment。三、前置条件Requirements根据总入口文档 .test-infra/kafka/strimzi/README.md使用这套 kustomization 前需要满足可连接的 Kubernetes 集群本仓库在.test-infra下提供了基于 Terraform 的 GKE 创建方式可参考 .test-infra/terraform/google-cloud-platform/google-kubernetes-engine该目录位于仓库根目录.test-infra/terraform/google-cloud-platform/google-kubernetes-engine预置集群kubectl CLI用于执行kubectl apply -k、kubectl get、kubectl port-forward等操作。四、预览 kustomizationPreview在真正向集群提交任何资源之前可以先本地渲染出 kustomization 的最终 YAML用于检查装配结果、确认资源清单是否符合预期kubectl kustomize folder例如预览01-strimzi-operatorkubectl kustomize 01-strimzi-operator执行后会在终端输出经namespace: strimzi归一化、并按依赖顺序合并好的全部 Operator 资源。这是一种低成本、无副作用的校验手段适合在 CI 或人工操作前先行核对版本目录与资源引用是否正确。五、第一步安装 Strimzi Operator总入口文档强调了一个关键纪律每一步命令必须等前一步完全完成后才能继续。安装 Operator 的命令为kubectl apply -k 01-strimzi-operator随后必须等待部署就绪再进入下一步创建 Kafka 集群kubectl get deploy strimzi-cluster-operator --namespace strimzi -w-w参数让命令持续监听 Deployment 状态直到strimzi-cluster-operator的可用副本数达到期望值后再Ctrl-C退出并继续。Operator 的组件构成manifest 拆解这一份 kustomization 装配的 Operator 由四类组件组成全部位于 v0.33.2 目录1. 服务账号ServiceAccount010-ServiceAccount-strimzi-cluster-operator.yaml定义了名为strimzi-cluster-operator的服务账号Operator 的 Pod 以它身份运行见 Deployment 中的serviceAccountName: strimzi-cluster-operator。2. RBAC 权限020/021/022/023-*一系列文件为 Operator 授予管理 Kafka 相关资源所需的集群级与命名空间级权限ClusterRole、RoleBinding、ClusterRoleBinding030/031/033-*则是 Operator 将权限委托给 Kafka broker、Entity Operator、Kafka 客户端时使用的委托绑定。正是这些 RBAC 让 Operator 能够监听kafka.strimzi.io组下的自定义资源并驱动实际 Kafka 集群的创建与滚动更新。3. CRD自定义资源定义040-Crd-kafka.yaml至049-Crd-kafkarebalance.yaml注册了 10 种 CRD包括核心的Kafka、以及KafkaTopic、KafkaUser、KafkaConnect、KafkaMirrorMaker、KafkaBridge、KafkaConnector、KafkaMirrorMaker2、KafkaRebalance、StrimziPodSet。安装后kubectl get crd可以看到这些资源类型后续步骤中的Kafka集群声明正是以 CRD 形式提交的。4. 配置与部署本体050-ConfigMap-strimzi-cluster-operator.yaml 提供 Operator 自身的log4j2.properties日志配置根日志级别可通过环境变量STRIMZI_LOG_LEVEL覆盖默认 INFO同时对 Kafka AdminClientorg.apache.kafka与 Zookeeperorg.apache.zookeeper单独降为 WARN 以抑制噪音Netty 与 OkHttp 保持 INFO。060-Deployment-strimzi-cluster-operator.yaml 是 Operator 本体单副本 Deployment镜像quay.io/strimzi/operator:0.33.2启动入口为/opt/strimzi/bin/cluster_operator_run.sh。Deployment 的关键环境变量从 060-Deployment-strimzi-cluster-operator.yaml 可以看到 Operator 运行时的核心调优项理解它们有助于排查后续 Kafka 集群创建问题环境变量值作用STRIMZI_NAMESPACE取自metadata.namespace即strimziOperator 监听与管理的命名空间STRIMZI_FULL_RECONCILIATION_INTERVAL_MS120000全量协调周期120 秒定期校正集群实际状态STRIMZI_OPERATION_TIMEOUT_MS300000单次操作超时300 秒超过即视为失败重试STRIMZI_KAFKA_IMAGES3.2.0...3.4.0...Kafka 版本到镜像的映射表创建集群时据此选择 broker 镜像STRIMZI_KAFKA_CONNECT_IMAGES/STRIMZI_KAFKA_MIRROR_MAKER_IMAGES/STRIMZI_KAFKA_MIRROR_MAKER_2_IMAGES同上分别对应 Kafka Connect 与两种 MirrorMaker 的版本映射STRIMZI_DEFAULT_TOPIC_OPERATOR_IMAGE/STRIMZI_DEFAULT_USER_OPERATOR_IMAGE/STRIMZI_DEFAULT_KAFKA_INIT_IMAGEquay.io/strimzi/operator:0.33.2Entity Operator 内 Topic/User Operator 及 broker 初始化容器的默认镜像STRIMZI_DEFAULT_KAFKA_BRIDGE_IMAGE/STRIMZI_DEFAULT_JMXTRANS_IMAGE/STRIMZI_DEFAULT_KANIKO_EXECUTOR_IMAGE/STRIMZI_DEFAULT_MAVEN_BUILDER0.33.2 系列镜像Bridge、JMX 转译、Kaniko 执行器与 Maven 构建器默认镜像STRIMZI_FEATURE_GATES空特性开关当前未启用任何实验特性STRIMZI_LEADER_ELECTION_ENABLEDtrue开启基于 Lease 的领导选举租约名strimzi-cluster-operator容器还配置了 HTTP 存活探针/healthy与就绪探针/ready初始延迟 10 秒、每 30 秒探测一次资源请求为 200m CPU / 384Mi 内存上限为 1000m CPU / 384Mi 内存并挂载了内存型emptyDir1Mi作为/tmp与 ConfigMap 挂载的/opt/strimzi/custom-config/。六、第二步创建 Kafka 持久化集群Operator 就绪后即可创建 Kafka 集群。这一步采用了 Kustomize 的base overlay模式基线定义集群规格overlay 针对特定环境打补丁。Base集群规格基线base/v0.33.2/kafka-persistent.yaml 定义了一个名为beam-testing-cluster的持久化 Kafka 集群该文件改编自 Strimzi 官方kafka-persistent-single.yaml示例apiVersion: kafka.strimzi.io/v1beta2 kind: Kafka metadata: name: beam-testing-cluster spec: kafka: version: 3.4.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: 100Gi deleteClaim: false zookeeper: replicas: 3 storage: type: persistent-claim size: 100Gi deleteClaim: false entityOperator: topicOperator: {} userOperator: {}要点解读3 个 broker、Kafka 3.4.0与 Operator v0.33.2 内置的STRIMZI_KAFKA_IMAGES映射表中的3.4.0条目一一对应关键 topic 配置了与副本数一致的复制因子offsets.topic.replication.factor3、transaction.state.log.replication.factor3、default.replication.factor3min.insync.replicas2保证高可用写入这是测试环境对数据可靠性的一种保守设定存储采用JBODtype: jbod 100Gi persistent-claimbroker 与 Zookeeper 各 3 副本deleteClaim: false表示集群删除时保留 PVC避免误删数据启用Entity OperatorTopic Operator User Operator支持通过KafkaTopic、KafkaUserCRD 管理主题与用户。OverlayGKE 内部负载均衡overlays/gke-internal-load-balanced/kustomization.yaml 在 base 之上追加补丁namespace: strimzi resources: - ../../base/v0.33.2 patchesStrategicMerge: - listeners.yaml补丁文件 listeners.yaml 为集群声明了三个监听器plain9092internal无 TLS集群内部客户端使用tls9093internalTLS集群内部启用 TLS 的客户端使用external9094type: loadbalancer无 TLS对外暴露bootstrap 与每个 broker 都通过注解cloud.google.com/load-balancer-type: Internal强制创建GKE 内部负载均衡器使集群仅在内网可达不向公网暴露 Kafka。应用该 overlay 的命令为kubectl apply -k 02-kafka-persistent/overlays/gke-internal-load-balanced可实时观察资源创建进度kubectl get all --namespace strimzi七、验证 Kafka 连接本地机器端口转发由于 Kafka 通过 GKE 内部负载均衡器暴露本地机器需要先做端口转发才能访问。文档给出了一条一次性完成的命令先按 selector 定位 broker Pod再把 Pod 的 9094 端口转发到本地kubectl port-forward --namespace strimzi $(kubectl get pod --namespace strimzi --selectorstrimzi.io/clusterbeam-testing-cluster,strimzi.io/kindKafka,strimzi.io/namebeam-testing-cluster-kafka --output jsonpath{.items[0].metadata.name}) 9094:9094该命令通过三个 Strimzi 标签精确锁定beam-testing-cluster的任一 Kafka broker Podstrimzi.io/clusterbeam-testing-cluster、strimzi.io/kindKafka、strimzi.io/namebeam-testing-cluster-kafka。简单连通性测试telnet在新终端中执行curl -v telnet://localhost:9094成功连接时会看到 TCP 层已建立* Trying 127.0.0.1:9094... * Connected to localhost (127.0.0.1) port 9094 (#0)这说明 9094 监听器链路Pod → 内部负载均衡 → port-forward已打通。使用 kcat 查看 broker 元数据进一步验证可以借助kcatKafka 命令行客户端安装方式见其官方仓库列出集群元数据kcat -L -b localhost:9094期望输出类似Metadata for all topics (from broker -1: localhost:9094/bootstrap): 3 brokers: broker 0 at 10.128.0.12:9094 (controller) broker 2 at 10.128.0.13:9094 broker 1 at 10.128.0.11:9094输出中 3 个 broker 的地址来自 GKE 内部负载均衡为每个 broker 分配的节点 IP证明集群具备完整的 3 副本拓扑。八、接入 Beam KafkaIODataflow 场景对于运行在 Dataflow 上的 Beam 管道无法依赖kubectl port-forward而应直接使用 Kafka 的对外 bootstrap 地址。先查询外部 bootstrap Servicekubectl get svc beam-testing-cluster-kafka-external-bootstrap --namespace strimzi典型输出NAME TYPE CLUSTER-IP EXTERNAL-IP PORT(S) AGE beam-testing-cluster-kafka-external-bootstrap LoadBalancer 10.167.128.247 10.128.0.14 9094:31331/TCP 86m取EXTERNAL-IP与端口9094拼接即得到 bootstrap server10.128.0.14:9094。随后在 KafkaIO 管道中使用KafkaIO.read().withBootstrapServers(10.128.0.14:9094)写入方向同理KafkaIO.write().withBootstrapServers(10.128.0.14:9094)withBootstrapServers是 KafkaIO 读取/写入配置的入口之一定义于 KafkaIO.java该方法接收逗号分隔的host:port列表Beam 管道据此与 Kafka 集群建立连接、消费或生产消息。整个strimzi目录的最终目的就是为这类 KafkaIO 集成测试持续提供稳定、可复现的 Kafka 服务。九、完整操作序列速查将上文命令按顺序汇总每步等待完成后再执行下一步# 1. 安装 Strimzi Operatorv0.33.2命名空间 strimzi kubectl apply -k 01-strimzi-operator kubectl get deploy strimzi-cluster-operator --namespace strimzi -w # 2. 创建 GKE 内部负载均衡的持久化 Kafka 集群 kubectl apply -k 02-kafka-persistent/overlays/gke-internal-load-balanced kubectl get all --namespace strimzi # 3. 本地验证端口转发 telnet kcat kubectl port-forward --namespace strimzi $(kubectl get pod --namespace strimzi --selectorstrimzi.io/clusterbeam-testing-cluster,strimzi.io/kindKafka,strimzi.io/namebeam-testing-cluster-kafka --output jsonpath{.items[0].metadata.name}) 9094:9094 curl -v telnet://localhost:9094 kcat -L -b localhost:9094 # 4. 供 Dataflow/KafkaIO 使用获取外部 bootstrap 地址 kubectl get svc beam-testing-cluster-kafka-external-bootstrap --namespace strimzi结语01-strimzi-operator是 Beam 测试基础设施中 Kafka 环境的基石它以 Kustomize 声明式装配 Strimzi v0.33.2 Operator 的全部组件命名空间、服务账号、RBAC、CRD、日志 ConfigMap 与 Deployment并以版本化目录的方式与后续02-kafka-persistent的集群基线保持同步。掌握这份 kustomization 的组装逻辑、版本对应关系与运维要点既能帮助你在本地复现 Beam 的 Kafka 集成测试环境也能为在 GKE 中自助搭建 Strimzi Kafka 提供一套可直接借鉴的参考实现。若需查看总入口要求与完整用法可继续阅读 .test-infra/kafka/strimzi/README.md。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam 测试基础设施使用 Strimzi 与 Kustomize 在 GKE 上部署持久化 Kafka 集群Apache Beam 测试基础设施使用 Strimzi 与 Kustomize 在 GKE 上部署持久化 Kafka 集群 导读 Apache Beam 的大数据批处理流处理数据工程Apache Beam 测试基础设施中的 Strimzi Kafka Operator 部署指南Terraform Helm Kustomize 全流程Apache Beam 测试基础设施中的 Strimzi Kafka Operator 部署指南Terraform Helm Kustomize 全流大数据批处理流处理数据工程Apache Beam 测试基础设施中的 Kafka 测试集群基于 Kubernetes 的 33 高可用部署实践Apache Beam 测试基础设施中的 Kafka 测试集群基于 Kubernetes 的 33 高可用部署实践 本文以 Apache Beam 仓库中大数据批处理流处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进