ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kubeless 复用已有 Kafka 集群:接入 Kafka Trigger、SASL 与 TLS 完整指南

Kubeless 复用已有 Kafka 集群:接入 Kafka Trigger、SASL 与 TLS 完整指南 后端云原生微服务【免费下载链接】kubelessKubernetes Native Serverless Framework项目地址https://gitcode.com/gh_mirrors/ku/kubeless点击查看免费下载导读本指南面向已经拥有自建 Kafka 集群、希望直接在 Kubeless 上启用 PubSub 函数的用户。文章将逐步演示如何通过自定义清单部署 Kafka 触发器控制器与 KafkaTrigger CRD并通过KAFKA_BROKERS、KAFKA_ENABLE_SASL、KAFKA_ENABLE_TLS等环境变量将 Kubeless 无缝接入已有集群最后给出端到端的消息收发验证方法。读完本文你将掌握在同一个 Kubernetes 集群内复用第三方 Kafka、配置认证加密并用kubeless topic/kubeless trigger kafka完成主题与触发器的日常管理。背景两种 Kafka 部署方式与本指南的适用场景Kubeless 的发布物中默认附带一组 Kafka 与 Zookeeper 的 StatefulSet 清单用于让用户快速拉起一套 开箱即用 的 PubSub 环境。这组 StatefulSet 会被部署在kubeless命名空间中并与 Kubeless 控制器、Kafka 触发器控制器一起工作。常规的部署方式参见 PubSub 事件说明 中的 Kafka 章节通过kubectl create -f应用官方提供的kafka-zookeeper-release.yaml即可。但很多团队实际上已经在同一个 Kubernetes 集群中运行着自己的 Kafka 集群。此时再部署一份 Kafka/Zookeeper 无疑浪费资源也会带来运维上的双重负担。本指南要解决的正是这一场景复用已有 Kafka 集群仅额外部署 Kubeless 的 Kafka 消费端kafka-trigger-controller与触发器自定义资源KafkaTrigger CRD让 PubSub 函数直接消费已有集群中的主题。需要明确的是Kafka 仅对 PubSub 类函数是必需的。如果你的函数只通过 HTTP 触发完全不需要关心本文内容参见 HTTP 触发器。前置条件确认已有 Kafka 与 Kubeless 环境确认 Kafka 集群运行状态假设你已有的 Kafka 集群运行在pubsub命名空间中且包含一个 Kafka brokerkafka-0和一个 Zookeeper 节点zoo-0$ kubectl -n pubsub get po NAME READY STATUS RESTARTS AGE kafka-0 1/1 Running 0 7h zoo-0 1/1 Running 0 7h $ kubectl -n pubsub get svc NAME TYPE CLUSTER-IP EXTERNAL-IP PORT(S) AGE kafka ClusterIP 10.55.253.151 none 9092/TCP 7h zookeeper ClusterIP 10.55.248.146 none 2181/TCP 7h其中kafkaService 的地址kafka.pubsub:9092将作为后续KAFKA_BROKERS环境变量的取值Kubernetes 集群内通过service.namespace形式访问。确认 Kubeless 控制面运行状态同时Kubeless 本体应已部署在kubeless命名空间中且控制器处于 Running 状态$ kubectl -n kubeless get po NAME READY STATUS RESTARTS AGE kubeless-controller-manager-58676964bb-l79gh 1/1 Running 0 5d关键前提为已有 Kafka 打上kubelesskafka标签文档中有一条容易被忽略的提示如果你希望使用kubeless topic命令需要为已有 Kafka 的 Pod 打上kubelesskafka标签CLI 才能定位到它。这条要求并非随意设计而是由 CLI 源码直接决定的。以 主题创建命令 为例createTopic会构造bash /opt/bitnami/kafka/bin/kafka-topics.sh --zookeeper zookeeper.namespace:2181 ...命令并通过execCommand在 Kafka Pod 内执行pods, err : utils.GetPodsByLabel(k8sClientSet, ctlNamespace, kubeless, kafka) if err ! nil { return fmt.Errorf(Cant find the kafka pod: %v, err) } else if len(pods.Items) 0 { return fmt.Errorf(Cant find any kafka pod) }即 CLI 通过标签选择器kubelesskafka查找目标 Pod详见 topicCreate.go找不到就会直接报错。主题发布命令 topicPublish.go 同样依赖该标签定位 broker 容器。因此请先为你的 Kafka 部署补充标签$ kubectl -n pubsub label pod kafka-0 kubelesskafka另外kubeless topic系列命令还支持通过--kafka-namespace指定 Kafka 控制器所在的命名空间默认值为kubeless见 topic.go 中的 Flag 定义。由于你的 Kafka 位于pubsub命名空间后续使用kubeless topic时需要显式传入$ kubeless topic create my-topic --kafka-namespace pubsub部署 Kafka 消费者与 KafkaTrigger CRD思路从通用清单中抽取三块内容Kubeless 提供的 Kafka 清单里除了 StatefulSet还包含了 Kafka 触发器控制器Deployment、kafkatriggers.kubeless.io自定义资源定义CRD以及配套的 RBAC 资源ClusterRole / ClusterRoleBinding。复用已有集群时只需把这些与控制面相关的资源抽取出来应用跳过 Kafka/Zookeeper StatefulSet 部分。最关键的一步是给触发器控制器容器添加环境变量KAFKA_BROKERS将其指向你已有的 Kafka 地址。这样 Kafka 触发器控制器才知道该从哪里消费消息。完整清单Deployment CRD RBAC将以下内容通过kubectl create -f -应用到集群输出即为各资源创建成功的回执$ echo --- apiVersion: apps/v1beta1 kind: Deployment metadata: labels: kubeless: kafka-trigger-controller name: kafka-trigger-controller namespace: kubeless spec: selector: matchLabels: kubeless: kafka-trigger-controller template: metadata: labels: kubeless: kafka-trigger-controller spec: containers: - image: bitnami/kafka-trigger-controller:latest imagePullPolicy: IfNotPresent name: kafka-trigger-controller env: - name: KAFKA_BROKERS value: kafka.pubsub:9092 # CHANGE THIS! serviceAccountName: controller-acct --- apiVersion: apiextensions.k8s.io/v1beta1 kind: CustomResourceDefinition metadata: name: kafkatriggers.kubeless.io spec: group: kubeless.io names: kind: KafkaTrigger plural: kafkatriggers singular: kafkatrigger scope: Namespaced version: v1beta1 --- apiVersion: rbac.authorization.k8s.io/v1beta1 kind: ClusterRoleBinding metadata: name: kafka-controller-deployer roleRef: apiGroup: rbac.authorization.k8s.io kind: ClusterRole name: kafka-controller-deployer subjects: - kind: ServiceAccount name: controller-acct namespace: kubeless --- apiVersion: rbac.authorization.k8s.io/v1beta1 kind: ClusterRole metadata: name: kafka-controller-deployer rules: - apiGroups: - resources: - services - configmaps verbs: - get - list - apiGroups: - kubeless.io resources: - functions - kafkatriggers verbs: - get - list - watch - update - delete | kubectl create -f - deployment kafka-trigger-controller created clusterrolebinding kafka-controller-deployer created clusterrole kafka-controller-deployer created customresourcedefinition kafkatriggers.kubeless.io created各资源的作用Deploymentkafka-trigger-controllerKafka 触发器控制器负责监听KafkaTriggerCRD 对象把函数与主题绑定并消费主题消息投递给对应函数。serviceAccountName: controller-acct表明它以该 ServiceAccount 身份运行。CRDkafkatriggers.kubeless.io定义KafkaTrigger资源group 为kubeless.ioversion 为v1beta1作用域为 Namespaced。Kubeless CLI 的kubeless trigger kafka create命令正是创建这种资源对象见 create.go 中构造kafkatriggers.kubeless.io/v1beta1KafkaTrigger 的代码。ClusterRole / ClusterRoleBindingkafka-controller-deployer授权触发器控制器读取services、configmaps并对kubeless.io组下的functions与kafkatriggers执行 get/list/watch/update/delete。这是控制器 Watch 函数与触发器、维护消费关系所需的权限边界。核心配置KAFKA_BROKERS 环境变量KAFKA_BROKERS是复用已有集群时唯一必须修改的配置项取值格式为host:port。本文示例中 Kafka Service 位于pubsub命名空间因此使用集群内部 DNS 名kafka.pubsub:9092。请根据你自己的 Service 名称与端口替换该值清单中已标注# CHANGE THIS!。若你的 Kafka 暴露了多个 broker可用逗号分隔的列表形式提供。端到端验证创建主题、发布消息、查看函数日志前提部署一个 PubSub 函数在验证之前需要先有一个被 Kafka 触发的函数。最简单的方式是部署一个打印event[data]的 Python 函数示例函数写法可参考 pubsub-functions.md 中的foobar完整的 Python 函数样例位于 examples/python$ kubeless function deploy pubsub-python --runtime python2.7 \ --handler test.foobar \ --from-file test.py然后使用kubeless trigger kafka将函数与主题关联。kubeless trigger kafka命令组包含 create/delete/list/update 四个子命令见 kafka_trigger.go其中 create 的--trigger-topic与--function-selector均为必填参数见 create.go$ kubeless trigger kafka create test --function-selector created-bykubeless,functionpubsub-python --trigger-topic s3-python这里--function-selector通过标签选择器定位函数 DeploymentKubeless 部署的函数默认带created-bykubeless标签--trigger-topic指定要监听的 Kafka 主题。在已有 Kafka 上创建主题由于触发器控制器不会替你创建主题你需要借助 Kafka 自带的脚本完成。本例直接复用 kafka Pod 内捆绑的二进制# create s3-python topic $ kubectl -n pubsub exec -it kafka-0 -- /opt/bitnami/kafka/bin/kafka-topics.sh --create --zookeeper zookeeper.pubsub:2181 --replication-factor 1 --partitions 1 --topic s3-python参数说明--zookeeper指定 Zookeeper 地址用于元数据管理--replication-factor 1与--partitions 1定义了副本数与分区数--topic s3-python为主题名。如果你的 Kafka 集群配置了更高副本/分区请按需调整。发布测试消息同样通过 kafka Pod 内的控制台生产者向主题写入消息$ kubectl -n pubsub exec -it kafka-0 -- /opt/bitnami/kafka/bin/kafka-console-producer.sh --broker-list localhost:9092 --topic s3-python hello world输入hello world并回车即可完成一次发布。验证函数消费另开一个终端实时查看 PubSub 函数的 Pod 日志应当能看到刚发布的消息被函数消费并打印$ kubectl logs -f pubsub-python-5445bdcb64-48bv2 hello world替代方案用kubeless topic管理主题如果不想手工执行kubectl exec也可以使用kubeless topic命令组create/list/delete/publish。它的底层实现正是kubectl exec到带kubelesskafka标签的 Kafka Pod 中执行kafka-topics.sh/kafka-console-producer.sh参见 topicCreate.go 与 topicPublish.go。由于你的 Kafka 在pubsub命名空间记得带上--kafka-namespace pubsub$ kubeless topic create s3-python --kafka-namespace pubsub Created topic s3-python. $ kubeless topic publish --topic s3-python --data hello world --kafka-namespace pubsubkubeless topic的详细用法create/list/delete 主题、发布消息等参见 pubsub-functions.md 的 Other commands 一节。启用 SASL 认证当你的已有 Kafka 集群启用了 SASL 认证时需要在触发器控制器的 Deployment 中追加三个环境变量KAFKA_ENABLE_SASL设为true开启 SASL 认证KAFKA_USERNAMESASL 用户名KAFKA_PASSWORDSASL 密码。考虑到安全实践文档特别提示密码类信息might use a secret建议用 Kubernetes Secret 承载敏感值再通过env.valueFrom.secretKeyRef注入而不是直接明文写在清单中。修改后的 Deployment 环境变量段如下$ echo --- apiVersion: apps/v1beta1 kind: Deployment metadata: labels: kubeless: kafka-trigger-controller name: kafka-trigger-controller namespace: kubeless spec: selector: matchLabels: kubeless: kafka-trigger-controller template: metadata: labels: kubeless: kafka-trigger-controller spec: containers: - image: bitnami/kafka-trigger-controller:latest imagePullPolicy: IfNotPresent name: kafka-trigger-controller env: ... - name: KAFKA_ENABLE_SASL value: true # CHANGE THIS! - name: KAFKA_USERNAME value: kafka-sasl-username # CHANGE THIS! - name: KAFKA_PASSWORD value: kafka-sasl-password # CHANGE THIS! ...修改后需要重新应用清单并滚动更新kafka-trigger-controllerDeployment使新环境变量生效。启用 TLS 加密通信当 Kafka 通信由 TLS 保护时必须设置KAFKA_ENABLE_TLS并按需指定以下变量KAFKA_CACERTSCA 证书路径用于校验服务器证书KAFKA_CERT与KAFKA_KEY客户端证书与私钥路径用于客户端证书校验KAFKA_INSECURE设为 true 可跳过 TLS 校验仅在测试环境使用生产环境不建议。前提创建保存证书与密钥的 Secret先将 CA 证书、客户端证书与私钥写入 Kubernetes Secret例如名为certs-and-keys-secret的 Secret包含ca.crt、cert.pem、key.pem$ kubectl -n kubeless create secret generic certs-and-keys-secret \ --from-fileca.crt./ca.crt \ --from-filecert.pem./cert.pem \ --from-filekey.pem./key.pem完整 TLS 示例清单--- apiVersion: apps/v1beta1 kind: Deployment metadata: labels: kubeless: kafka-trigger-controller name: kafka-trigger-controller namespace: kubeless spec: selector: matchLabels: kubeless: kafka-trigger-controller template: metadata: labels: kubeless: kafka-trigger-controller spec: volumes: - name: kafka-volume secret: secretName: certs-and-keys-secret # REPLACE WITH SECRET HOLDING CERTS AND KEYS containers: - image: bitnami/kafka-trigger-controller:latest imagePullPolicy: IfNotPresent name: kafka-trigger-controller volumeMounts: - name: kafka-volume mountPath: /path/to/certsandkeys env: ... - name: KAFKA_ENABLE_TLS value: true # ENABLE TLS - name: KAFKA_CACERTS value: /path/to/certsandkeys/ca.crt # CHANGE THIS! (NOTE : PATH HERE MATCHING THE MOUNT PATH ABOVE) - name: KAFKA_CERT value: /path/to/certsandkeys/cert.pem # CHANGE THIS! (NOTE : PATH HERE MATCHING THE MOUNT PATH ABOVE) - name: KAFKA_KEY value: /path/to/certsandkeys/key.pem # CHANGE THIS! (NOTE : PATH HERE MATCHING THE MOUNT PATH ABOVE) ...清单中三个证书路径必须与volumeMounts.mountPath示例为/path/to/certsandkeys保持一致——环境变量里的路径是容器内文件系统路径它由 Secret 卷挂载产生因此要确保KAFKA_CACERTS、KAFKA_CERT、KAFKA_KEY指向的正是挂载点下的文件名。相关命令与常见问题KafkaTrigger 的日常管理除 create 外kubeless trigger kafka还提供kubeless trigger kafka list列出当前命名空间的 Kafka 触发器kubeless trigger kafka delete trigger-name删除触发器kubeless trigger kafka update trigger-name --trigger-topic new-topic --function-selector labels修改触发器绑定的主题或函数标签选择器源码实现见 update.go未传入的参数不会被改动。常见问题kubeless topic报 Cant find any kafka pod说明集群中不存在带kubelesskafka标签的 Pod请为已有 Kafka Pod 打标签并确认使用了正确的--kafka-namespace。kafka-trigger-controller 反复重启优先检查KAFKA_BROKERS地址是否正确可达kafka.pubsub:9092是否与你的 Service 名/端口一致、SASL/TLS 变量与证书路径是否匹配。函数收不到消息确认主题已创建、--trigger-topic与发布所用主题完全一致且--function-selector能命中已部署函数的标签。关于存储的提醒Kubeless 默认附带的 Kafka StatefulSet 使用 PVC在部分平台上需要预先创建 PV 或配置动态存储否则 Kafka Pod 会一直 Pending——不过这正是本文通过复用已有集群要规避的问题。相关排查思路可参考 Troubleshooting 文档。延伸阅读PubSub 事件机制Kafka / NATS 触发器的整体工作方式与默认部署清单Triggers 总览Kubeless 全部触发器类型索引实现新的触发器如需扩展自定义触发器可以参考的规范。赞分享后端云原生微服务【免费下载链接】kubelessKubernetes Native Serverless Framework项目地址https://gitcode.com/gh_mirrors/ku/kubeless点击查看免费下载相关推荐kafka-go SASL/SCRAM认证完整指南安全连接Kafka集群的终极教程kafka go SASL/SCRAM认证完整指南安全连接Kafka集群的终极教程 在当今数据驱动的时代Apache Kafka已成为企业级数据流处理的核心消息队列后端数据工程Kafka-UI开源Apache Kafka集群管理工具完整指南Kafka UI开源Apache Kafka集群管理工具完整指南 项目概述 Kafka UI是一款免费开源的Web界面工具专门用于管理和监控Apache K后端前端监控大盘终极指南kafka-go TLS加密与SASL认证实战教程终极指南kafka go TLS加密与SASL认证实战教程 在当今数据安全日益重要的环境下保护Kafka消息传输的安全性变得至关重要。kafka go作为G消息队列后端数据工程上一篇Koodo Reader 上手六平台免费开源电子书阅读器十分钟配好书库同步下一篇CANN/ge获取附属流IDAPI创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表