# Apache Kafka สำหรับวิศวกรข้อมูล: Streaming, Partitions และคำถามสัมภาษณ์ > เจาะลึก Apache Kafka สำหรับวิศวกรข้อมูล ครอบคลุมสถาปัตยกรรม streaming กลยุทธ์ partition consumer groups และคำถามสัมภาษณ์ที่พบบ่อย พร้อมตัวอย่างการใช้งานจริงด้วย Kafka 4.x และ KRaft - Published: 2026-04-20 - Updated: 2026-04-26 - Author: SharpSkill - Tags: kafka, streaming, data-engineering, partitions, consumer-groups, kraft, event-driven - Reading time: 11 min --- Apache Kafka เป็นหัวใจของ stack ด้าน data engineering สมัยใหม่เกือบทุกแห่ง รองรับเหตุการณ์หลายล้านล้านครั้งต่อวันในองค์กรทุกขนาด ด้วย Kafka 4.x ที่ทำงานบน KRaft อย่างเต็มรูปแบบ (ไม่ต้องพึ่ง ZooKeeper อีกต่อไป) แพลตฟอร์มนี้จึงดูแลรักษาง่ายขึ้นในขณะที่ยังคงการรับประกันแบบเดิมที่ทำให้ Kafka กลายเป็นมาตรฐานอุตสาหกรรมสำหรับ data pipeline แบบเรียลไทม์ > **Kafka 4.x: KRaft เท่านั้น** > > ตั้งแต่ Apache Kafka 4.0 เป็นต้นไป ไม่จำเป็นต้องใช้ ZooKeeper อีกต่อไป KRaft (Kafka Raft) จัดการ metadata ทั้งหมดแบบ native ลดความซับซ้อนในการดำเนินงานและตัดส่วนประกอบโครงสร้างพื้นฐานออกจาก deployment ระดับ production ## สถาปัตยกรรม Streaming ของ Kafka สำหรับ Data Pipeline Kafka ทำงานเหมือน distributed commit log ผู้ผลิต (producer) เพิ่ม record ต่อท้าย partition ของ topic และผู้บริโภค (consumer) อ่าน record เหล่านั้นตามลำดับ การออกแบบแบบ append-only นี้รองรับทั้ง streaming แบบเรียลไทม์และการ replay แบบ batch จาก offset ใดก็ได้ ทำให้ Kafka เหมาะสมเป็นพิเศษกับงาน data engineering ที่ต้องการรูปแบบทั้งสองอย่าง คลัสเตอร์ Kafka ประกอบด้วย broker (เซิร์ฟเวอร์) topic (ช่องทางเชิงตรรกะ) และ partition (shard เชิงกายภาพของ topic) แต่ละ partition คือลำดับ record ที่เรียงตามลำดับและไม่สามารถเปลี่ยนแปลงได้ broker หนึ่งตัวทำหน้าที่เป็น partition leader จัดการการอ่านและเขียนทั้งหมด ในขณะที่ follower replica เก็บสำเนาไว้เพื่อ fault tolerance คุณสมบัติทางสถาปัตยกรรมที่สำคัญ: **การเรียงลำดับถูกรับประกันภายใน partition ไม่ใช่ข้าม partition** ข้อจำกัดเดียวนี้กำหนดการตัดสินใจในการออกแบบ Kafka ส่วนใหญ่ใน data pipeline ```yaml # docker-compose.yml — Kafka 4.x KRaft cluster (no ZooKeeper) services: kafka-1: image: apache/kafka:4.2.0 environment: KAFKA_NODE_ID: 1 KAFKA_PROCESS_ROLES: broker,controller KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093 KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_LOG_DIRS: /var/lib/kafka/data KAFKA_NUM_PARTITIONS: 6 KAFKA_DEFAULT_REPLICATION_FACTOR: 3 KAFKA_MIN_INSYNC_REPLICAS: 2 ports: - "9092:9092" ``` การตั้งค่านี้เริ่ม broker โหมด KRaft ที่ทำหน้าที่เป็น controller ด้วย ค่า `KAFKA_CONTROLLER_QUORUM_VOTERS` มาแทนที่สิ่งที่ ZooKeeper เคยจัดการ ได้แก่ การเลือก leader, metadata ของ topic และสมาชิกของคลัสเตอร์ ## กลยุทธ์ Partition ที่กำหนดประสิทธิภาพของ Pipeline Partition กำหนดทั้ง parallelism การเรียงลำดับ และ throughput การเลือกจำนวน partition และกลยุทธ์ key ที่เหมาะสมส่งผลโดยตรงต่อความน่าเชื่อถือของ pipeline และความสามารถในการขยายขนาดของ consumer **แนวทางการเลือกจำนวน partition:** - เริ่มต้นด้วยจำนวน consumer ที่คาดว่าจะประมวลผล topic พร้อมกัน - แต่ละ partition รองรับ throughput การเขียนได้ประมาณ 10 MB/วินาทีบนฮาร์ดแวร์สมัยใหม่ - จำนวน partition ที่มากขึ้นจะเพิ่ม end-to-end latency เล็กน้อย (มีการเลือก leader มากขึ้น มี file handle มากขึ้น) - Kafka 4.x จัดการ partition หลายพันรายการต่อ broker ได้อย่างมีประสิทธิภาพ ด้วยการปรับปรุง metadata ของ KRaft การเลือก partition key กำหนดว่า record ใดจะลงไปยัง partition ใด record ที่มี key เดียวกันจะไปยัง partition เดียวกันเสมอ คงการเรียงลำดับสำหรับ key นั้นไว้ ```python # partition_strategy.py — Choosing the right partition key from confluent_kafka import Producer import json def create_producer(): return Producer({ 'bootstrap.servers': 'kafka-1:9092', 'acks': 'all', 'enable.idempotence': True, 'max.in.flight.requests.per.connection': 5 }) def publish_order_event(producer, event): # Key by customer_id: all events for one customer stay ordered key = str(event['customer_id']).encode('utf-8') value = json.dumps(event).encode('utf-8') producer.produce( topic='order-events', key=key, value=value, headers=[('source', b'order-service'), ('version', b'2')] ) producer.flush() def publish_clickstream(producer, event): # Key by session_id: ordering within a session matters # NOT by user_id — one user may have multiple sessions key = str(event['session_id']).encode('utf-8') value = json.dumps(event).encode('utf-8') producer.produce( topic='clickstream', key=key, value=value ) ``` การเลือก `customer_id` แทน `session_id` เป็น partition key สะท้อนความต้องการการเรียงลำดับที่แตกต่างกัน เหตุการณ์คำสั่งซื้อต้องการการเรียงลำดับต่อลูกค้าเพื่อรักษาความสอดคล้องของธุรกรรม เหตุการณ์ clickstream ต้องการการเรียงลำดับต่อ session เพื่อสร้างเส้นทางผู้ใช้ขึ้นมาใหม่ได้อย่างถูกต้อง > **ความเสี่ยงจาก hot partition** > > หาก key หนึ่งสร้าง record มากกว่า key อื่นอย่างมีนัยสำคัญ (เช่น ลูกค้าระดับองค์กรเพียงรายเดียวสร้างเหตุการณ์ 80%) partition นั้นจะกลายเป็นคอขวด ควรติดตามตัวชี้วัด partition lag และพิจารณาใช้ composite key หรือ custom partitioner สำหรับ workload ที่กระจายไม่สมดุล ## Consumer Groups และการจัดการ Offset Consumer group ช่วยให้ประมวลผล topic แบบขนานได้ Kafka จัดสรรแต่ละ partition ให้กับ consumer หนึ่งตัวภายในกลุ่มเท่านั้น ดังนั้น parallelism สูงสุดจึงเท่ากับจำนวน partition การจัดการ offset ควบคุมความหมายของการส่งมอบแบบ exactly-once เทียบกับ at-least-once Kafka เก็บ offset ที่ commit ไว้ในหัวข้อภายในชื่อ `__consumer_offsets` ช่วงเวลาของการ commit offset เทียบกับการประมวลผลกำหนดการรับประกันการส่งมอบ ```python # consumer_pipeline.py — Consumer group with manual offset management from confluent_kafka import Consumer, KafkaError import json def create_consumer(group_id: str): return Consumer({ 'bootstrap.servers': 'kafka-1:9092', 'group.id': group_id, 'auto.offset.reset': 'earliest', 'enable.auto.commit': False, # Manual commit for at-least-once 'max.poll.interval.ms': 300000, 'session.timeout.ms': 45000, 'isolation.level': 'read_committed' # Only read committed transactions }) def process_batch(consumer: Consumer, batch_size: int = 500): """Process records in micro-batches for throughput.""" consumer.subscribe(['order-events']) buffer = [] while True: msg = consumer.poll(timeout=1.0) if msg is None: continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: continue raise Exception(msg.error()) record = json.loads(msg.value().decode('utf-8')) buffer.append(record) if len(buffer) >= batch_size: # Process the full batch (write to warehouse, etc.) write_to_warehouse(buffer) # Commit AFTER successful processing consumer.commit(asynchronous=False) buffer.clear() def write_to_warehouse(records: list): # Insert batch into data warehouse (BigQuery, Snowflake, etc.) pass ``` การตั้งค่า `enable.auto.commit` เป็น `False` และ commit หลังจาก `write_to_warehouse` สำเร็จเท่านั้นรับประกันการส่งมอบแบบ at-least-once หาก consumer ขัดข้องหลังการประมวลผลแต่ก่อน commit record จะถูกส่งซ้ำเมื่อเริ่มใหม่ สำหรับ data pipeline ที่ป้อนข้อมูลเข้าสู่ warehouse นี่เป็นการแลกเปลี่ยนที่ถูกต้องโดยทั่วไป การเขียนแบบ idempotent เข้าสู่ warehouse จัดการรายการที่ซ้ำกันได้ ## Kafka Connect สำหรับการเชื่อมต่อ Data Pipeline Kafka Connect จัดเตรียม framework ที่ขยายขนาดได้สำหรับการเคลื่อนย้ายข้อมูลระหว่าง Kafka กับระบบภายนอกโดยไม่ต้องเขียนโค้ด producer/consumer แบบกำหนดเอง source connector นำข้อมูลเข้าสู่ Kafka ส่วน sink connector ผลักข้อมูลออกไปข้างนอก ```json { "name": "postgres-cdc-source", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.hostname": "postgres-primary", "database.port": "5432", "database.user": "replication_user", "database.dbname": "production", "topic.prefix": "cdc", "table.include.list": "public.orders,public.customers,public.products", "slot.name": "debezium_slot", "plugin.name": "pgoutput", "publication.name": "dbz_publication", "snapshot.mode": "initial", "transforms": "route", "transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.route.regex": "cdc\\.public\\.(.*)", "transforms.route.replacement": "warehouse.$1" } } ``` Debezium CDC connector นี้จับทุกการ INSERT, UPDATE และ DELETE จากตาราง PostgreSQL และเผยแพร่เป็นเหตุการณ์ Kafka การแปลง `RegexRouter` เปลี่ยนชื่อหัวข้อจาก `cdc.public.orders` เป็น `warehouse.orders` รักษา namespace ของ pipeline ให้สะอาด เมื่อรวมกับ BigQuery หรือ Snowflake sink connector จะสร้าง CDC pipeline ที่ทำงานอัตโนมัติเต็มรูปแบบโดยไม่ต้องใช้โค้ดแบบกำหนดเอง ## ความหมายแบบ Exactly-Once ใน Kafka Pipeline Kafka รองรับความหมายแบบ exactly-once (EOS) ผ่าน idempotent producer และ transactional API สำหรับ data engineering pipeline EOS ป้องกันการมี record ซ้ำในระบบปลายทางโดยไม่ต้องใช้ logic การ deduplicate องค์ประกอบสามอย่างที่ทำให้ EOS ทำงานได้: 1. **Idempotent producer**: producer แต่ละตัวได้รับ Producer ID ที่ไม่ซ้ำกัน broker จะกำจัดการลองซ้ำที่ซ้ำซ้อนโดยใช้ sequence number ดังนั้นการลองส่งซ้ำผ่านเครือข่ายจึงไม่สร้างรายการที่ซ้ำกัน 2. **Transaction**: producer สามารถเขียนเข้าหลาย partition และ commit offset ในธุรกรรมเดียวแบบ atomic การเขียนทุกครั้งสำเร็จหรือไม่มีอะไรสำเร็จเลย 3. **การแยกแบบ `read_committed`**: consumer ที่ตั้งค่า `isolation.level=read_committed` จะเห็นเฉพาะ record จากธุรกรรมที่ commit แล้วเท่านั้น ```python # exactly_once_pipeline.py — Transactional consume-transform-produce from confluent_kafka import Consumer, Producer import json def transactional_pipeline(): consumer = Consumer({ 'bootstrap.servers': 'kafka-1:9092', 'group.id': 'etl-transformer', 'enable.auto.commit': False, 'isolation.level': 'read_committed' }) producer = Producer({ 'bootstrap.servers': 'kafka-1:9092', 'transactional.id': 'etl-transformer-001', 'acks': 'all', 'enable.idempotence': True }) producer.init_transactions() consumer.subscribe(['raw-events']) while True: msg = consumer.poll(1.0) if msg is None or msg.error(): continue # Begin atomic transaction producer.begin_transaction() try: raw = json.loads(msg.value()) enriched = transform(raw) # Write transformed record to output topic producer.produce( 'enriched-events', key=msg.key(), value=json.dumps(enriched).encode('utf-8') ) # Commit consumer offset within the same transaction producer.send_offsets_to_transaction( consumer.position(consumer.assignment()), consumer.consumer_group_metadata() ) producer.commit_transaction() except Exception: producer.abort_transaction() def transform(record: dict) -> dict: # Apply business logic transformations record['processed_at'] = '2026-04-20T00:00:00Z' record['pipeline_version'] = '2.1' return record ``` ค่า `transactional.id` ช่วยให้ broker สามารถ fence zombie producer ได้ หาก consumer ขัดข้องและ instance ใหม่เริ่มต้นด้วย `transactional.id` เดียวกัน broker จะยกเลิกธุรกรรมที่ค้างอยู่จาก instance เก่า ป้องกัน output ที่ซ้ำกัน ## Share Groups: ความหมายของคิวใน Kafka 4.2 Kafka 4.2 แนะนำ Share Groups ที่พร้อมใช้งานระดับ production นำความหมายของคิวข้อความแบบดั้งเดิมมาสู่ Kafka ต่างจาก consumer group ที่แต่ละ partition ถูกจัดสรรให้กับ consumer หนึ่งตัวเท่านั้น Share Groups อนุญาตให้ consumer หลายตัวประมวลผล record จาก partition เดียวกันได้อย่างอิสระ | คุณสมบัติ | Consumer Groups | Share Groups | |---|---|---| | การจัดสรร partition | เฉพาะตัว (consumer หนึ่งตัวต่อ partition) | แชร์ (consumer หลายตัวต่อ partition) | | การรับประกันการเรียงลำดับ | คงการเรียงลำดับต่อ partition | ไม่มีการรับประกันการเรียงลำดับ | | กรณีใช้งาน | สตรีมเหตุการณ์ที่เรียงลำดับ | การกระจายงาน คิวงาน | | การยืนยันรับ | อิงตาม offset (commit position) | ต่อ record (acknowledge แยกกัน) | | Parallelism สูงสุด | จำกัดด้วยจำนวน partition | จำกัดด้วยจำนวน consumer | Share Groups แก้ข้อจำกัดที่มีมานาน: การขยายจำนวน consumer ให้เกินจำนวน partition สำหรับ workload เช่น การส่งการแจ้งเตือนหรือการกระจายงาน batch ที่การเรียงลำดับไม่สำคัญ Share Groups ทำให้ไม่จำเป็นต้องมีระบบคิวภายนอกควบคู่ไปกับ Kafka > **เมื่อใดควรใช้ Share Groups** > > Share Groups เหมาะกับ workload การกระจายงาน เช่น การส่งอีเมล การประมวลผลการอัปโหลด การรัน ML inference สำหรับ event sourcing, CDC หรือ pipeline ใดที่ต้องการการประมวลผลแบบเรียงลำดับ consumer group ยังคงเป็นตัวเลือกที่ถูกต้อง ## คำถามสัมภาษณ์ Kafka สำหรับวิศวกรข้อมูล การสัมภาษณ์ทางเทคนิคสำหรับตำแหน่ง data engineering มักทดสอบความรู้ Kafka คำถามต่อไปนี้ครอบคลุมแนวคิดที่ถูกประเมินบ่อยที่สุด **Q: Kafka รับประกันการเรียงลำดับข้อความอย่างไร?** การเรียงลำดับถูกรับประกันเฉพาะภายใน partition เดียวเท่านั้น record ทั้งหมดที่มี partition key เดียวกันจะไปยัง partition เดียวกันและถูกเพิ่มต่อท้ายตามลำดับ ระหว่าง partition ไม่มีการเรียงลำดับใดอยู่ การเลือก partition key ที่ถูกต้องเป็นกลไกหลักในการควบคุมการเรียงลำดับใน Kafka pipeline **Q: เกิดอะไรขึ้นเมื่อ consumer ในกลุ่มล้มเหลว?** Group coordinator ตรวจพบความล้มเหลวผ่าน heartbeat ที่หายไป (ควบคุมโดย `session.timeout.ms`) จะกระตุ้น rebalance จัดสรร partition ของ consumer ที่ล้มเหลวใหม่ให้กับสมาชิกที่เหลือในกลุ่ม ในระหว่าง rebalance การบริโภคจะหยุดชั่วคราว Cooperative sticky rebalancing (ค่าเริ่มต้นใน Kafka 4.x) ลดการหยุดชะงักโดยจัดสรรใหม่เฉพาะ partition ที่ได้รับผลกระทบ **Q: อธิบายความแตกต่างระหว่าง `acks=1` และ `acks=all`** ด้วย `acks=1` broker ยืนยันการเขียนหลังจาก leader replica เก็บ record ไว้ ด้วย `acks=all` broker จะรอจนกว่า in-sync replica (ISR) ทั้งหมดจะยืนยันการเขียน `acks=all` รวมกับ `min.insync.replicas=2` รับประกันว่าจะไม่สูญหายข้อมูลหาก broker เดี่ยวล้มเหลว แลกกับ latency ที่สูงขึ้นเล็กน้อย **Q: คุณจะออกแบบ CDC pipeline โดยใช้ Kafka อย่างไร?** จับการเปลี่ยนแปลงจาก source database โดยใช้ Debezium (สำหรับ PostgreSQL, MySQL) หรือ CDC connector แบบ native เผยแพร่เหตุการณ์การเปลี่ยนแปลงไปยัง topic Kafka ที่ partition ตาม primary key ใช้ sink connector (BigQuery Sink, S3 Sink) เพื่อโหลดการเปลี่ยนแปลงเข้าสู่ analytical warehouse ตั้งค่า `isolation.level=read_committed` บน consumer เพื่อหลีกเลี่ยงการอ่านธุรกรรม database ที่ยังไม่ commit ใช้ schema registry เพื่อจัดการ schema Avro หรือ Protobuf สำหรับเหตุการณ์การเปลี่ยนแปลง **Q: ISR คืออะไรและทำไมจึงสำคัญ?** ISR (In-Sync Replicas) คือชุดของ replica ที่ตามทัน partition leader อย่างเต็มที่ มีเพียงสมาชิก ISR เท่านั้นที่มีสิทธิ์ได้รับเลือกเป็น leader ค่า `min.insync.replicas` กำหนดว่ามี replica กี่ตัวต้องยืนยันการเขียนก่อนที่จะถือว่า commit แล้ว หาก ISR ลดลงต่ำกว่า `min.insync.replicas` broker จะปฏิเสธการเขียนแทนที่จะเสี่ยงต่อการสูญหายข้อมูล สำหรับการเตรียมตัวสัมภาษณ์ data engineering เพิ่มเติม ชุด [คำถามสัมภาษณ์ data engineering](/th/technologies/data-engineering) ครอบคลุมหัวข้อที่กว้างขึ้นรวมถึง [รูปแบบ pipeline ETL/ELT](/th/technologies/data-engineering/interview-questions/etl-elt-patterns) และ [การสร้างแบบจำลองข้อมูล](/th/technologies/data-engineering/interview-questions/data-modeling-de) ## บทสรุป - Kafka 4.x กับ KRaft กำจัด ZooKeeper อย่างสมบูรณ์ ลดภาระการดำเนินงานสำหรับทีม data engineering - การเลือก partition key กำหนดการรับประกันการเรียงลำดับ ควรเลือก key ตามความต้องการของการประมวลผลปลายทาง ไม่ใช่ตามความสะดวก - การ commit offset ด้วยตนเองด้วย `enable.auto.commit=False` รองรับการส่งมอบแบบ at-least-once transactional API ให้ความหมายแบบ exactly-once สำหรับ pipeline แบบ consume-transform-produce - Kafka Connect กับ Debezium ให้ CDC ที่พร้อมใช้งานระดับ production โดยไม่ต้องเขียนโค้ด consumer แบบกำหนดเอง - Share Groups (Kafka 4.2) เพิ่มความหมายของคิวสำหรับ workload การกระจายงานที่ไม่ต้องการการเรียงลำดับ - การขยายขนาด consumer group ถูกจำกัดด้วยจำนวน partition ควรวางแผนจำนวน partition ตามความต้องการ parallelism ของ consumer ที่จุดสูงสุด --- Source: SharpSkill (https://sharpskill.dev), tech interview preparation for your real stack. HTML version of this page: https://sharpskill.dev/th/blog/data-engineering/apache-kafka-data-engineering-streaming-partitions