当前位置: 首页 > news >正文

K8s 部署 Kafka (KRaft) + SASL/SCRAM-SHA-512 踩坑与终极实战指南

这是一份基于前面排坑与实践沉淀的Kafka (KRaft 模式) + SASL/SCRAM-SHA-512 安全认证的完整 Helm 部署教程。

架构包含了声明式的用户管理、动态注册脚本、全流程对齐的 SCRAM 加密机制以及高可用存储配置。

📖 教程目录

  1. 项目目录结构

  2. 完整配置文件

    • values.yaml

    • templates/configmap.yaml

    • templates/secret.yaml

    • templates/statefulset.yaml

    • templates/service.yaml

  3. 部署与生命周期管理

  4. 验证与客户端接入

1. 项目目录结构

在 Helm Chart 根目录下按如下结构组织文件:

Plaintext

agent-intent-kafka/ ├── Chart.yaml ├── values.yaml └── templates/ ├── configmap.yaml ├── secret.yaml ├── service.yaml └── statefulset.yaml

2. 完整配置文件

📄values.yaml

用于集中配置集群镜像、端口、存储以及 SASL 用户密码。

YAML

kafka: enabled: true replicaCount: 1 clusterId: "agent-intent-kafka-cluster-id" image: repository: apache/kafka tag: 3.7.0 pullPolicy: IfNotPresent ports: plain: 9092 sasl: 9094 controller: 9093 storage: 10Gi storageClass: "" # 根据实际集群填写,为空则使用默认 StorageClass resources: limits: cpu: "2" memory: 4Gi requests: cpu: "500m" memory: 2Gi auth: enabled: true users: admin: "admin-secret-pass" consumer: "Consumer@CKD6UDVah" producer: "Producer@CKD6UDVah"

📄templates/configmap.yaml

定义 Kafka KRaft 核心服务配置(包含角色分配、监听器与 StandardAuthorizer 鉴权类)。

YAML

{{- if .Values.kafka.enabled }} apiVersion: v1 kind: ConfigMap metadata: name: {{ include "agent-intent.fullname" . }}-kafka-config labels: {{- include "agent-intent.labels" . | nindent 4 }} app.kubernetes.io/component: kafka data: server.properties: | # --- KRaft 核心集群角色定义 --- process.roles=broker,controller controller.listener.names=CONTROLLER early.start.listeners=CONTROLLER # --- 认证与授权配置 --- {{- if .Values.kafka.auth.enabled }} listeners=PLAINTEXT://0.0.0.0:{{ .Values.kafka.ports.plain }},SASL_PLAINTEXT://0.0.0.0:{{ .Values.kafka.ports.sasl }},CONTROLLER://0.0.0.0:{{ .Values.kafka.ports.controller }} listener.security.protocol.map=PLAINTEXT:PLAINTEXT,SASL_PLAINTEXT:SASL_PLAINTEXT,CONTROLLER:PLAINTEXT sasl.enabled.mechanisms=SCRAM-SHA-512 sasl.mechanism.inter.broker.protocol=PLAINTEXT authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer super.users=User:admin;User:ANONYMOUS allow.everyone.if.no.acl.found=false {{- else }} listeners=PLAINTEXT://0.0.0.0:{{ .Values.kafka.ports.plain }},CONTROLLER://0.0.0.0:{{ .Values.kafka.ports.controller }} listener.security.protocol.map=PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT {{- end }} # --- 存储与数据目录 --- log.dirs=/var/lib/kafka/data # --- 副本因子配置 --- {{- $replicaCount := .Values.kafka.replicaCount | int }} {{- $rf := $replicaCount }} {{- $minIsr := $replicaCount | int }} {{- if gt $minIsr 1 }}{{ $minIsr = sub $minIsr 1 }}{{ end }} {{- if lt $minIsr 1 }}{{ $minIsr = 1 }}{{ end }} default.replication.factor={{ $rf }} offsets.topic.replication.factor={{ $rf }} transaction.state.log.replication.factor={{ $rf }} transaction.state.log.min.isr={{ $minIsr }} min.insync.replicas={{ $minIsr }} # --- 性能调优 --- num.io.threads=16 num.network.threads=8 num.partitions=3 log.retention.hours=168 {{- end }}

📄templates/secret.yaml

集中定义服务端/客户端 JAAS 配置文件,严格排除注释与特殊字符,对齐ScramLoginModule

YAML

{{- if .Values.kafka.enabled }} apiVersion: v1 kind: Secret metadata: name: {{ include "agent-intent.fullname" . }}-kafka-auth labels: {{- include "agent-intent.labels" . | nindent 4 }} app.kubernetes.io/component: kafka type: Opaque stringData: cluster-id: "{{ .Values.kafka.clusterId | default (randAlphaNum 16) }}" {{- if .Values.kafka.auth.enabled }} {{- $adminPassword := index .Values.kafka.auth.users "admin" | default "admin-secret-pass" }} kafka_server_jaas.conf: | KafkaServer { org.apache.kafka.common.security.scram.ScramLoginModule required username="admin" password="{{ $adminPassword }}"; }; Client { org.apache.kafka.common.security.scram.ScramLoginModule required username="admin" password="{{ $adminPassword }}"; }; client.properties: | security.protocol=SASL_PLAINTEXT sasl.mechanism=SCRAM-SHA-512 sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="admin" password="{{ $adminPassword }}"; {{- end }} {{- end }}

📄templates/statefulset.yaml

核心控制器,内置(...) &后台异步注册任务,在服务连通瞬间自动将 SCRAM 凭证写入 KRaft 元数据。

YAML

{{- if .Values.kafka.enabled }} apiVersion: apps/v1 kind: StatefulSet metadata: name: {{ include "agent-intent.fullname" . }}-kafka labels: {{- include "agent-intent.labels" . | nindent 4 }} app.kubernetes.io/component: kafka spec: serviceName: {{ include "agent-intent.fullname" . }}-kafka-headless replicas: {{ .Values.kafka.replicaCount }} podManagementPolicy: Parallel selector: matchLabels: {{- include "agent-intent.selectorLabels" . | nindent 6 }} app.kubernetes.io/component: kafka template: metadata: labels: {{- include "agent-intent.selectorLabels" . | nindent 8 }} app.kubernetes.io/component: kafka spec: containers: - name: kafka image: "{{ .Values.kafka.image.repository }}:{{ .Values.kafka.image.tag }}" imagePullPolicy: {{ .Values.kafka.image.pullPolicy }} command: - /bin/bash - -ec - | ID=${HOSTNAME##*-} CONFIG=/tmp/server.properties # 复制只读模板到 /tmp 生成动态配置文件 cp /etc/kafka/server.properties.template ${CONFIG} echo "node.id=${ID}" >> ${CONFIG} {{- $replicaCount := .Values.kafka.replicaCount | int }} {{- $fullname := include "agent-intent.fullname" . }} {{- $namespace := .Release.Namespace }} {{- $controllerPort := .Values.kafka.ports.controller }} {{- if eq $replicaCount 1 }} echo "controller.quorum.voters=0@localhost:{{ $controllerPort }}" >> ${CONFIG} {{- else }} echo "controller.quorum.voters={{- range $i := until $replicaCount }}{{ if gt $i 0 }},{{ end }}{{ $i }}@{{ $fullname }}-kafka-{{ $i }}.{{ $fullname }}-kafka-headless.{{ $namespace }}.svc.cluster.local:{{ $controllerPort }}{{- end }}" >> ${CONFIG} {{- end }} {{- if .Values.kafka.auth.enabled }} export ADVERTISE_HOST="${HOSTNAME}.{{ include "agent-intent.fullname" . }}-kafka-headless.{{ .Release.Namespace }}.svc.cluster.local" echo "advertised.listeners=PLAINTEXT://${ADVERTISE_HOST}:{{ .Values.kafka.ports.plain | default 9092 }},SASL_PLAINTEXT://${ADVERTISE_HOST}:{{ .Values.kafka.ports.sasl | default 9094 }}" >> ${CONFIG} {{- else }} export ADVERTISE_HOST="${HOSTNAME}.{{ include "agent-intent.fullname" . }}-kafka-headless.{{ .Release.Namespace }}.svc.cluster.local" echo "advertised.listeners=PLAINTEXT://${ADVERTISE_HOST}:{{ .Values.kafka.ports.plain | default 9092 }}" >> ${CONFIG} {{- end }} if [ ! -f /var/lib/kafka/data/meta.properties ]; then /opt/kafka/bin/kafka-storage.sh format \ --ignore-formatted \ -t "${CLUSTER_ID}" \ -c ${CONFIG} fi {{- if .Values.kafka.auth.enabled }} # 后台安全异步注入 SCRAM 凭证(轮询本地 9092 端口,连通即注入,不阻塞主流程) ( while ! </dev/tcp/127.0.0.1/{{ .Values.kafka.ports.plain | default 9092 }}; do sleep 0.5 done {{- $adminPassword := index .Values.kafka.auth.users "admin" | default "admin-secret-pass" }} /opt/kafka/bin/kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \ --entity-type users --entity-name admin --alter \ --add-config 'SCRAM-SHA-512=[password={{ $adminPassword }}]' || true {{- range $user, $pass := .Values.kafka.auth.users }} {{- if ne $user "admin" }} /opt/kafka/bin/kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \ --entity-type users --entity-name {{ $user }} --alter \ --add-config 'SCRAM-SHA-512=[password={{ $pass }}]' || true {{- end }} {{- end }} ) & {{- end }} exec /opt/kafka/bin/kafka-server-start.sh ${CONFIG} ports: - containerPort: {{ .Values.kafka.ports.plain | default 9092 }} name: plain-port {{- if .Values.kafka.auth.enabled }} - containerPort: {{ .Values.kafka.ports.sasl | default 9094 }} name: sasl-port {{- end }} startupProbe: exec: command: - /bin/bash - -c - "</dev/tcp/127.0.0.1/{{ .Values.kafka.ports.plain | default 9092 }}" initialDelaySeconds: 1 periodSeconds: 2 timeoutSeconds: 2 failureThreshold: 30 livenessProbe: exec: command: - /bin/bash - -c - "</dev/tcp/127.0.0.1/{{ .Values.kafka.ports.plain | default 9092 }}" periodSeconds: 15 timeoutSeconds: 5 failureThreshold: 3 readinessProbe: exec: command: - /bin/bash - -c - "</dev/tcp/127.0.0.1/{{ .Values.kafka.ports.plain | default 9092 }}" initialDelaySeconds: 2 periodSeconds: 5 timeoutSeconds: 3 failureThreshold: 3 env: - name: POD_NAMESPACE valueFrom: fieldRef: fieldPath: metadata.namespace - name: CLUSTER_ID valueFrom: secretKeyRef: name: {{ include "agent-intent.fullname" . }}-kafka-auth key: cluster-id {{- if .Values.kafka.auth.enabled }} - name: KAFKA_OPTS value: "-Djava.security.auth.login.config=/etc/kafka/jaas/kafka_server_jaas.conf" {{- end }} - name: KAFKA_HEAP_OPTS value: "-Xmx2G -Xms2G" volumeMounts: - name: kafka-data mountPath: /var/lib/kafka/data - name: kafka-config mountPath: /etc/kafka/server.properties.template subPath: server.properties readOnly: true {{- if .Values.kafka.auth.enabled }} - name: kafka-jaas-config mountPath: /etc/kafka/jaas/kafka_server_jaas.conf subPath: kafka_server_jaas.conf readOnly: true - name: kafka-jaas-config mountPath: /etc/kafka/client.properties subPath: client.properties readOnly: true {{- end }} resources: {{- toYaml .Values.kafka.resources | nindent 10 }} volumes: - name: kafka-config configMap: name: {{ include "agent-intent.fullname" . }}-kafka-config {{- if .Values.kafka.auth.enabled }} - name: kafka-jaas-config secret: secretName: {{ include "agent-intent.fullname" . }}-kafka-auth {{- end }} volumeClaimTemplates: - metadata: name: kafka-data spec: accessModes: [ "ReadWriteOnce" ] resources: requests: storage: {{ .Values.kafka.storage }} {{- if .Values.kafka.storageClass }} storageClassName: {{ .Values.kafka.storageClass }} {{- end }} {{- end }}

📄templates/service.yaml

为外部暴露集群 Headless 服务,内部与外部认证端口隔离。

YAML

{{- if .Values.kafka.enabled }} apiVersion: v1 kind: Service metadata: name: {{ include "agent-intent.fullname" . }}-kafka-headless namespace: {{ .Release.Namespace }} labels: {{- include "agent-intent.labels" . | nindent 4 }} app.kubernetes.io/component: kafka app.kubernetes.io/part-of: kafka spec: clusterIP: None publishNotReadyAddresses: true ports: - name: plaintext port: {{ .Values.kafka.ports.plain | default 9092 }} targetPort: {{ .Values.kafka.ports.plain | default 9092 }} protocol: TCP {{- if .Values.kafka.auth.enabled }} - name: sasl-plaintext port: {{ .Values.kafka.ports.sasl | default 9094 }} targetPort: {{ .Values.kafka.ports.sasl | default 9094 }} protocol: TCP {{- end }} - name: controller port: {{ .Values.kafka.ports.controller | default 9093 }} targetPort: {{ .Values.kafka.ports.controller | default 9093 }} protocol: TCP selector: {{- include "agent-intent.selectorLabels" . | nindent 4 }} app.kubernetes.io/component: kafka {{- end }}

3. 部署与生命周期管理

部署 Chart

使用 Helm 将组件安装至命名空间:

Bash

helm upgrade --install agent-intent-kafka ./agent-intent-kafka \ -n agent-intent-system \ --create-namespace

查看运行状态

Bash

kubectl get pods -n agent-intent-system -l app.kubernetes.io/component=kafka

4. 验证与客户端接入

🔐 命令行测试(容器内)

1. 查询 Topic 列表(验证 SASL_PLAINTEXT 握手与鉴权)

Bash

kubectl exec -it agent-intent-kafka-0 -n agent-intent-system -- \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server 127.0.0.1:9094 \ --command-config /etc/kafka/client.properties \ --list

2. 创建 Topic 与收发消息验证

Bash

# 创建 Topic kubectl exec -it agent-intent-kafka-0 -n agent-intent-system -- \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server 127.0.0.1:9094 \ --command-config /etc/kafka/client.properties \ --create --topic demo-topic --partitions 1 --replication-factor 1 # 写入测试消息 kubectl exec -it agent-intent-kafka-0 -n agent-intent-system -- \ bash -c 'echo "hello kafka scram sha 512" | /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server 127.0.0.1:9094 \ --producer.config /etc/kafka/client.properties \ --topic demo-topic' # 消费测试消息 kubectl exec -it agent-intent-kafka-0 -n agent-intent-system -- \ /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server 127.0.0.1:9094 \ --consumer.config /etc/kafka/client.properties \ --topic demo-topic \ --from-beginning --max-messages 1

💻 业务微服务连接配置参考

跨集群或同 Namespace 微服务接入时,请在应用服务配置中增加如下参数(以 Python / Java 为例):

  • Bootstrap Server:agent-intent-kafka-0.agent-intent-kafka-headless.agent-intent-system.svc.cluster.local:9094

  • Security Protocol:SASL_PLAINTEXT

  • SASL Mechanism:SCRAM-SHA-512

  • Username:admin

  • Password:admin-secret-pass

http://www.jsqmd.com/news/1247835/

相关文章:

  • 三星冰箱SAMSUNG推出全国统一24小时售后服务电话人工上线2026最新公布 - 优企名品
  • 【实战】黄金暴涨破 $4100 与 A 股深 V 大反弹!如何用 Python + QuantDash 实现跨市场(A股/美股/贵金属)联动套利监控
  • 2026年7月最新!雅典济南售后热线全网同步,网点地址一览,客户服务更贴心 - 亨得利官方服务中心
  • 提示工程架构师实验室的技术创新与实践
  • TRAE智能体开发:模块化构建与实战优化
  • TPS61183 WLED驱动芯片设计实战:从原理到PCB布局与调试
  • X-AnyLabeling 全平台保姆级安装教程(Windows / Linux / macOS)
  • 伯爵中国售后服务中心|地址及24小时客服电话权威信息声明(2026年7月更新) - 亨得利官方服务中心
  • 解决Unity Web Player更新失败:从原理到实战的完整指南
  • 升学季选比较好的美国签证办理课程中心 6个清单参考
  • AI视频角色崩坏?揭秘LLM+Diffusion跨帧ID锚定技术:从面部微表情到服装纹理的7层一致性校验协议
  • 低成本论文降AI方案:TextHumanizer与StyleTransferPro实战
  • 学术论文智能降重技术解析与应用实践
  • XUnity.AutoTranslator 游戏实时翻译插件:从原理到实战的深度优化指南
  • lu,生理药理实验多用仪
  • 企业如何合法解除不胜任员工劳动合同?
  • 行业规范同步公示,2026 上海合扬黄金回收全域网点服务细则一览 - 生活商业速报
  • GEO优化要沉淀什么?广拓时代谈AI时代的品牌认知资产
  • 深圳龙华葡萄牙语培训哪个好 - GrowUME
  • 物联网之安规下的 MQTT Payload 加密实践
  • 基于YOLOv5的智能火灾检测系统设计与优化
  • RAG应用开发实战:文本分块与向量检索核心技术解析
  • MIGM-Shortcut:AI图像生成4倍加速技术解析
  • 文件搜索慢到想砸电脑?这款神器秒搜全盘,比系统自带快N倍!
  • UI自动化之Playwright简介
  • 和平精英超体对抗新角色介绍 怎么用电脑控手机玩和平精英
  • 照着CSDN老教程敲代码,我成功把App Store审核搞黄了
  • 论文格式排版总被打回?[特殊字符]2026国标规范一键搞定,告别格式扣分
  • 2026年AI安全公司哪家强?全网靠谱品牌深度评测与选型攻略
  • 【选型指南】2026量化开发:akshare、Tushare 与 QuantDash 多维度客观对比与数据清洗实测