การรัน Kafka แบบ KRaft บน Kubernetes สรุปได้เป็นการตัดสินใจไม่กี่เรื่อง คือ controller จะรันที่ไหน, ใช้ StatefulSet ที่ broker แต่ละตัวมี persistent volume ของตัวเอง, มี headless Service ที่ให้ DNS ชื่อคงที่กับ broker ทุกตัว และตั้ง advertised.listeners เป็นชื่อนั้น เพื่อให้ client ไปถึง partition leader ได้ถูกตัว ถ้าทำครบ Kafka บน Kubernetes จะทำงานแทบไม่ต่างจากบน VM แต่ถ้าตั้ง advertised.listeners ผิด client จะต่อได้ ดึง metadata ได้ แล้วล้มทุก request หลังจากนั้น
เรื่อง topic, partition และ consumer group อธิบายไว้แล้วใน Apache Kafka partition และ consumer group บทความนี้เน้นการดูแล cluster ได้แก่บทบาทของแต่ละ node, KRaft, การตั้งค่า replication และ manifest บน Kubernetes
Broker กับ controller: data plane และ control plane ของ Kafka
Kafka cluster มีสองชั้น คล้ายกับที่ Kubernetes แยก control plane ออกจาก worker node
Broker (data plane) เก็บ partition ลง disk, รับการอ่านเขียนจาก client และ replicate ข้อมูลจาก leader ไปยัง follower broker แต่ละตัวถือทั้ง leader และ follower ของหลาย topic ปนกัน และ replica จะไม่อยู่บน broker เดียวกับ leader ของมัน Kafka พยายามกระจาย leader ให้เท่า ๆ กัน เพื่อไม่ให้โหลดการเขียนไปกองที่เครื่องเดียว
Controller (control plane) ไม่ได้เก็บข้อความของคุณ แต่เก็บ metadata ของ cluster เช่นมี topic อะไรบ้าง, broker ไหนเป็น leader ของ partition ไหน, replica อยู่ที่ใด และ broker ตัวไหนยังมีชีวิต เมื่อ broker ตาย controller ตัวที่ active จะเลือก leader ใหม่จาก replica ที่ in-sync อยู่
Controller รันเป็น quorum ปกติ 3 ตัว หรือ 5 ตัวใน cluster ใหญ่ มีตัวหนึ่ง active ส่วนที่เหลือตาม metadata log และพร้อมขึ้นแทน ต้องมีเสียงข้างมากทำงานอยู่ controller 3 ตัวจึงทนตายได้ 1 ตัว และ 5 ตัวทนได้ 2 ตัว
ZooKeeper vs KRaft
ก่อนมี KRaft ชั้น control plane คือ ZooKeeper ensemble ที่แยกออกมาต่างหาก ทำให้ต้อง deploy, ทำ security, monitor และ upgrade ระบบกระจายสองระบบพร้อมกัน
KRaft (Kafka Raft) ย้าย controller เข้ามาอยู่ใน Kafka เอง metadata เก็บเป็น log ที่ replicate ด้วย Raft ภายใน controller quorum
| ZooKeeper mode | KRaft mode | |
|---|---|---|
| ส่วนประกอบ | Kafka broker และ ZooKeeper ensemble | Kafka อย่างเดียว |
| Metadata | เก็บใน ZooKeeper | metadata log ที่ replicate ภายใน Kafka |
| Controller failover | controller ตัวใหม่ต้องโหลด state จาก ZooKeeper | controller สำรองมี log อยู่แล้ว |
| สถานะ | deprecated ตั้งแต่ 3.5, ถูกถอดออกใน 4.0 | พร้อมใช้ production ตั้งแต่ 3.3, เป็นโหมดเดียวตั้งแต่ 4.0 |
ประโยชน์ที่เห็นจริงคือดูแลระบบเดียว, controller failover เร็วขึ้น และรองรับ partition ต่อ cluster ได้มากขึ้นมาก cluster ใหม่ควรใช้ KRaft ส่วน cluster เดิมที่ใช้ ZooKeeper ต้อง migrate ไป KRaft บน release 3.x ก่อน จึงจะ upgrade เป็น 4.x ได้
Process role
KRaft node ทุกตัวต้องตั้ง process.roles:
brokerเก็บข้อมูลและให้บริการ clientcontrollerเป็นสมาชิกของ metadata quorumbroker,controllerคือ combined mode ตั้งค่าง่ายกว่าและพอใช้สำหรับ development หรือ cluster เล็ก แต่เอกสารของ Kafka ไม่แนะนำสำหรับ deployment ที่สำคัญ เพราะ broker ที่งานหนักอาจทำให้ controller ช้าลง
KRaft ใช้ node.id ไม่ใช่ broker.id แบบยุค ZooKeeper และ broker กับ controller ใช้ ID ชุดเดียวกัน
Checklist การตั้งค่าแบบหลาย broker
- มี broker อย่างน้อย 3 ตัว เพื่อให้ topic ใช้ replication factor 3 ได้
- ใช้ controller แยก 3 ตัว ใน production (5 ตัวสำหรับ cluster ใหญ่มาก)
-
default.replication.factor=3และตั้งค่าเดียวกันให้ internal topic คือoffsets.topic.replication.factorและtransaction.state.log.replication.factor -
min.insync.replicas=2โดย producer ใช้acks=all -
unclean.leader.election.enable=false(ค่าเริ่มต้น) เพื่อไม่ให้ replica ที่ข้อมูลตามไม่ทันขึ้นเป็น leader แล้วทำข้อมูลหายแบบเงียบ ๆ - Disk แยกของใครของมัน สำหรับ broker แต่ละตัว ห้ามแชร์
-
listenersและadvertised.listenersเป็น address ที่ client เข้าถึงได้จริง - ตั้ง rack awareness (
broker.rack) เป็น availability zone เพื่อให้ replica ของ partition เดียวกันอยู่คนละ zone
Replication factor 3 คู่กับ min.insync.replicas 2 ให้อะไรบ้าง
เมื่อตั้ง replication.factor=3 แต่ละ partition จะมี leader หนึ่งตัวและ follower สองตัวอยู่คนละ broker เมื่อใช้ min.insync.replicas=2 คู่กับ acks=all การเขียนจะสำเร็จก็ต่อเมื่อมี replica อย่างน้อยสองตัวได้รับข้อมูลแล้ว ใน cluster ที่มี 3 broker:
| Broker ที่ล่ม | ข้อมูล | การเขียนด้วย acks=all |
|---|---|---|
| 0 | ปลอดภัย | รับได้ |
| 1 | ปลอดภัย | รับได้ เหลือ in-sync replica สองตัว |
| 2 | ยังอยู่บน replica ตัวเดียว | ถูกปฏิเสธด้วย NotEnoughReplicas จนกว่า broker จะกลับมา |
Cluster ยังรับการเขียนได้เมื่อ broker ล่มหนึ่งตัว ซึ่งครอบคลุมกรณี rolling restart หรือ drain node และเลือกปฏิเสธการเขียนแทนที่จะเสี่ยงทำข้อมูลที่ ack แล้วหายเมื่อเจอเหตุใหญ่กว่านั้น การตั้ง min.insync.replicas เท่ากับ replication factor ฟังดูปลอดภัยกว่า แต่ผลคือแค่ restart broker ตัวเดียวก็เขียนไม่ได้แล้ว
Producer หา partition leader อย่างไร
Client ไม่ได้ส่งทุก record ผ่าน load balancer แต่ route เอง:
- Client ต่อไปที่ address ไหนก็ได้ใน
bootstrap.serversแล้วขอ metadata - Metadata บอก address ที่ broker แต่ละตัว advertise ไว้ และบอกว่า leader ของแต่ละ partition อยู่ที่ไหน
- สำหรับแต่ละ record producer เลือก partition จาก key, หา leader ของ partition นั้น แล้วส่ง batch ตรงไปที่ broker ตัวนั้น
- ถ้า leader ย้าย broker จะตอบ error เช่น
NOT_LEADER_OR_FOLLOWERclient ก็ refresh metadata แล้ว retry
นี่คือเหตุผลที่ advertised.listeners สำคัญมาก bootstrap address แค่พา client เข้าประตู หลังจากนั้น client จะต่อไปยัง address ที่ broker แต่ละตัว advertise ไว้ ถ้า broker advertise เป็น localhost:9092 หรือ pod IP ที่ client เข้าไม่ถึง bootstrap จะผ่าน แต่ produce request จะล้ม
ฝั่ง consumer ก็อ่านจาก partition leader เป็นค่าเริ่มต้นเหมือนกัน ถ้าตั้ง broker.rack ที่ broker, client.rack ที่ consumer และ replica.selector.class=org.apache.kafka.common.replica.RackAwareReplicaSelector consumer จะอ่านจาก follower ที่อยู่ zone เดียวกันได้ ช่วยลด traffic ข้าม zone
Kafka KRaft บน Kubernetes: ชิ้นส่วนหลัก
| ความต้องการ | Kubernetes object | เหตุผล |
|---|---|---|
| Identity คงที่ต่อ broker | StatefulSet | pod ชื่อ kafka-0, kafka-1, kafka-2 และใช้ชื่อเดิมทุกครั้งที่ restart |
| Disk ของตัวเองต่อ broker | volumeClaimTemplates | pod แต่ละตัวได้ PersistentVolumeClaim ของตัวเอง ย้าย node ก็ตามไปด้วย |
| Address ตรงถึง broker แต่ละตัว | Headless Service | DNS หนึ่งชื่อต่อ pod |
| จุดเชื่อมต่อแรก | ClusterIP Service | bootstrap.servers address เดียวที่ไปถึง broker ตัวไหนก็ได้ |
| ดูแลระบบได้อย่างปลอดภัย | PodDisruptionBudget | ล่มได้ไม่เกินหนึ่ง broker ระหว่าง voluntary disruption |
Deployment ไม่เหมาะกับงานนี้ เพราะ pod ได้ชื่อสุ่ม ไม่มี identity คงที่ และมี volume แยกต่อ pod ไม่ได้ เปรียบเทียบละเอียดอยู่ใน Kubernetes workload
Headless Service ตัวเดียวครอบคลุม broker ทุกตัว
ความเข้าใจผิดที่พบบ่อยคือคิดว่าต้องมี Service แยกให้ broker ทุกตัว จริง ๆ ไม่ต้อง headless Service (clusterIP: None) ไม่มี virtual IP และไม่ทำ load balancing เมื่อใช้คู่กับ StatefulSet มันจะสร้าง DNS record ให้ทุก pod เช่น kafka-0.kafka-headless.kafka.svc.cluster.local ถ้า scale เป็น 4 broker ตัว kafka-3 ก็ได้ชื่อของตัวเองอัตโนมัติ
เพิ่ม ClusterIP Service ปกติอีกตัวไว้สำหรับ bootstrap เพราะ broker ตัวไหนก็ตอบ metadata request ได้ ส่วน producer และ consumer ไม่ต้องมี Service ของตัวเอง เพราะเป็นฝ่ายต่อออกไปอย่างเดียว ไม่มีใครเรียกเข้ามาหา
apiVersion: v1
kind: Service
metadata:
name: kafka-headless
namespace: kafka
spec:
clusterIP: None
publishNotReadyAddresses: true # controllers must resolve each other before they are Ready
selector:
app: kafka
ports:
- name: broker
port: 9092
- name: controller
port: 9093
---
apiVersion: v1
kind: Service
metadata:
name: kafka-bootstrap
namespace: kafka
spec:
selector:
app: kafka
ports:
- name: broker
port: 9092StatefulSet ของ KRaft สาม node
Manifest นี้รัน node แบบ combined mode 3 ตัวด้วย image ทางการ apache/kafka เหมาะสำหรับทำความเข้าใจส่วนประกอบและใช้ใน development ได้ สำหรับ production ให้แยก controller กับ broker เป็นคนละ StatefulSet หรือใช้ operator
apiVersion: apps/v1
kind: StatefulSet
metadata:
name: kafka
namespace: kafka
spec:
serviceName: kafka-headless
replicas: 3
podManagementPolicy: Parallel
selector:
matchLabels:
app: kafka
template:
metadata:
labels:
app: kafka
spec:
securityContext:
fsGroup: 1000
containers:
- name: kafka
image: apache/kafka:3.9.1
ports:
- name: broker
containerPort: 9092
- name: controller
containerPort: 9093
env:
- name: POD_NAME
valueFrom:
fieldRef:
fieldPath: metadata.name
- name: KAFKA_NODE_ID
valueFrom:
fieldRef:
fieldPath: metadata.labels['apps.kubernetes.io/pod-index']
- name: CLUSTER_ID
value: "q1Sh-9_ISia_zwGINzRvyQ" # generate once: kafka-storage.sh random-uuid
- name: KAFKA_PROCESS_ROLES
value: "broker,controller"
- name: KAFKA_LISTENERS
value: "PLAINTEXT://:9092,CONTROLLER://:9093"
- name: KAFKA_ADVERTISED_LISTENERS
value: "PLAINTEXT://$(POD_NAME).kafka-headless.kafka.svc.cluster.local:9092"
- name: KAFKA_LISTENER_SECURITY_PROTOCOL_MAP
value: "PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT"
- name: KAFKA_CONTROLLER_LISTENER_NAMES
value: "CONTROLLER"
- name: KAFKA_INTER_BROKER_LISTENER_NAME
value: "PLAINTEXT"
- name: KAFKA_CONTROLLER_QUORUM_VOTERS
value: "[email protected]:9093,[email protected]:9093,[email protected]:9093"
- name: KAFKA_LOG_DIRS
value: "/var/lib/kafka/data/logs"
- name: KAFKA_DEFAULT_REPLICATION_FACTOR
value: "3"
- name: KAFKA_MIN_INSYNC_REPLICAS
value: "2"
- name: KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR
value: "3"
- name: KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR
value: "3"
- name: KAFKA_TRANSACTION_STATE_LOG_MIN_ISR
value: "2"
readinessProbe:
tcpSocket:
port: 9092
periodSeconds: 10
resources:
requests:
cpu: "1"
memory: 4Gi
volumeMounts:
- name: data
mountPath: /var/lib/kafka/data
volumeClaimTemplates:
- metadata:
name: data
spec:
accessModes: ["ReadWriteOnce"]
resources:
requests:
storage: 100Giรายละเอียดที่สำคัญ:
podManagementPolicy: Parallelสั่งให้ทุก pod เริ่มพร้อมกัน ถ้าใช้ค่าเริ่มต้นOrderedReadykafka-1จะรอจนกว่าkafka-0Ready แต่ controller quorum ต้องมี pod รันอยู่เป็นเสียงข้างมากก่อนถึงจะทำงานได้node.idมาจาก pod index labelapps.kubernetes.io/pod-index(Kubernetes 1.28 ขึ้นไป) ให้เลขลำดับของ pod ซึ่งตรงกับ ID ใน voter listadvertised.listenersใช้ DNS ของ pod ตัวเอง metadata จึงบอก client ได้ตรงว่า leader แต่ละตัวอยู่ที่ไหน- Log directory เป็น subdirectory ของ volume volume หลายแบบมีโฟลเดอร์
lost+foundอยู่ที่ root ซึ่ง Kafka จะพยายามโหลดเป็น partition CLUSTER_IDสร้างครั้งเดียวและห้ามเปลี่ยน image จะใช้ค่านี้ format storage ที่ยังว่างตอนเริ่มครั้งแรก
การตั้งค่านี้รองรับเฉพาะ client ภายใน cluster ถ้า client อยู่นอก Kubernetes ต้องเข้าถึง broker ได้ทีละตัว ซึ่งหมายถึงต้องมี LoadBalancer หรือ NodePort ต่อ broker และ listener ตัวที่สองที่ advertise address ภายนอกเหล่านั้น
เพิ่ม PodDisruptionBudget เพื่อให้การ drain node ทำทีละ broker แล้วตรวจ quorum และสร้าง topic:
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata:
name: kafka
namespace: kafka
spec:
maxUnavailable: 1
selector:
matchLabels:
app: kafkakubectl -n kafka exec kafka-0 -- /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server localhost:9092 describe --status
kubectl -n kafka exec kafka-0 -- /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 --create --topic bookings \
--partitions 12 --replication-factor 3รัน Kafka ด้วย operator: Strimzi
ถ้าเขียน manifest เอง คุณต้องรับผิดชอบ rolling restart ที่ดู ISR, การหมุน certificate, rack awareness, external listener และการ scale เอง operator คือการเขียนความรู้เหล่านี้ลงเป็นโค้ด Strimzi เป็น Kafka operator ที่ใช้กันแพร่หลายที่สุด คุณประกาศ resource Kafka และ KafkaNodePool สำหรับ controller และ broker แล้วมันจะสร้าง pod, volume และ Service ให้ Strimzi รุ่นใหม่ ๆ รองรับเฉพาะ KRaft แล้ว
Strimzi ใช้รูปแบบ Service สองตัวแบบเดียวกับข้างบน คือ bootstrap Service (<cluster>-kafka-bootstrap) และ headless Service สำหรับ broker (<cluster>-kafka-brokers) ไม่ว่าจะมี 3 หรือ 30 broker หลักการทำงานของ operator โดยทั่วไปอธิบายไว้ใน Kubernetes operator และ CRD
คำถามที่พบบ่อย
Kafka ยังต้องใช้ ZooKeeper อยู่ไหม
ไม่ต้องแล้ว KRaft พร้อมใช้ production ตั้งแต่ Kafka 3.3 และ Kafka 4.0 ถอดการรองรับ ZooKeeper ออกทั้งหมด
ต้องมี KRaft controller กี่ตัว
Cluster ส่วนใหญ่ใช้ 3 ตัว ทนตายได้ 1 ตัว ถ้า 5 ตัวทนได้ 2 ตัว ควรใช้จำนวนคี่ เพราะ controller ตัวที่ 4 ไม่ได้เพิ่มความทนทานขึ้นเลย
ทำไม client ต่อ Kafka ได้แต่ produce ไม่ได้
เกือบทุกครั้งคือ advertised.listeners bootstrap ผ่าน แต่ broker advertise address ที่ client เข้าไม่ถึง ให้ดูว่า broker แต่ละตัว advertise อะไรไว้ แล้วลอง resolve จาก network ของ client
Kafka บน Kubernetes ควรใช้ Deployment หรือ StatefulSet
StatefulSet หรือ operator ที่จัดการ pod โดยให้การรับประกันแบบเดียวกัน broker ต้องมีชื่อคงที่ ID คงที่ และ disk ของตัวเอง
สรุป
- ใช้ KRaft เพราะ ZooKeeper mode ถูกถอดออกแล้วใน Kafka 4.0
- ใช้ controller แยก 3 ตัวใน production และเก็บ combined mode ไว้ใช้ตอน development
- Replication factor 3,
min.insync.replicas=2และacks=allคือ baseline มาตรฐานด้าน durability - StatefulSet, volume ต่อ broker, headless Service และ bootstrap Service คือชิ้นส่วนหลักบน Kubernetes
- ตั้ง
advertised.listenersเป็น address ที่ client ทุกตัวเข้าถึงได้
Kafka บน Kubernetes ต่อยอดจาก StatefulSet, storage และ Service โดยตรง ถ้าอยากปูพื้นเรื่องเหล่านี้ให้แน่น คอร์ส DevOps ฟรีของ Vectorkub ครอบคลุมพื้นฐานชุดนี้
