Apache Kafka สำหรับวิศวกรข้อมูล: Streaming, Partitions และคำถามสัมภาษณ์

เจาะลึก Apache Kafka สำหรับวิศวกรข้อมูล ครอบคลุมสถาปัตยกรรม streaming กลยุทธ์ partition consumer groups และคำถามสัมภาษณ์ที่พบบ่อย พร้อมตัวอย่างการใช้งานจริงด้วย Kafka 4.x และ KRaft

แผนผังสถาปัตยกรรม streaming ของ Apache Kafka พร้อม partition และการไหลของข้อมูล

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 ที่ทำงานอัตโนมัติเต็มรูปแบบโดยไม่ต้องใช้โค้ดแบบกำหนดเอง

พร้อมที่จะพิชิตการสัมภาษณ์ Data Engineering แล้วหรือยังครับ?

ฝึกฝนด้วยตัวจำลองแบบโต้ตอบ, flashcards และแบบทดสอบเทคนิคครับ

ความหมายแบบ 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 GroupsShare 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 ครอบคลุมหัวข้อที่กว้างขึ้นรวมถึง รูปแบบ pipeline ETL/ELT และ การสร้างแบบจำลองข้อมูล

เริ่มฝึกซ้อมเลย!

ทดสอบความรู้ของคุณด้วยตัวจำลองสัมภาษณ์และแบบทดสอบเทคนิคครับ

บทสรุป

  • 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 ที่จุดสูงสุด

เริ่มฝึกซ้อมเลย!

ทดสอบความรู้ของคุณด้วยตัวจำลองสัมภาษณ์และแบบทดสอบเทคนิคครับ

ชาเลนจ์ประจำวัน

คุณหาบั๊กใน Data Engineering เจอไหม

โค้ดจริงหนึ่งชิ้น บั๊กที่ซ่อนอยู่หนึ่งจุด วันละหนึ่งครั้ง ลองได้โดยไม่ต้องมีบัญชี

Anthony Fillion-Maillet

เขียนโดย

Anthony Fillion-Maillet

ผู้ก่อตั้ง SharpSkill

เป็นนักพัฒนาฟูลสแตกมากว่า 10 ปี ดูแล SharpSkill และรับผิดชอบทุกสิ่งที่เผยแพร่ที่นี่

อัปเดตเมื่อ 26 เมษายน 2569

แท็ก

#kafka
#streaming
#data-engineering
#partitions
#consumer-groups
#kraft
#event-driven

แชร์

บทความที่เกี่ยวข้อง

Great Expectations data quality validation framework diagram

Great Expectations 2026: การตรวจสอบคุณภาพข้อมูลและคำถามสัมภาษณ์

คู่มือฉบับสมบูรณ์เกี่ยวกับ framework Great Expectations 1.22 สำหรับการตรวจสอบคุณภาพข้อมูลใน pipeline Python รวมถึงการผสานกับ Airflow และคำถามสัมภาษณ์ data engineering

Apache Beam vs Spark 2026 เปรียบเทียบสำหรับ data engineering

Apache Beam vs Spark 2026: เปรียบเทียบ Pipeline แบบรวมและคำถามสัมภาษณ์

คู่มือครบถ้วนเปรียบเทียบ Apache Beam 2.76 และ Spark 4.2 สำหรับ data engineering เรียนรู้ความแตกต่างด้านสถาปัตยกรรม windowing ประสิทธิภาพ และคำถามสัมภาษณ์ที่พบบ่อย

dbt data transformations testing interview 2026

dbt ในปี 2026: การแปลงข้อมูล การทดสอบ และคำถามสัมภาษณ์งาน

คู่มือ dbt สำหรับวิศวกรข้อมูล: การแปลง SQL, การสร้างโมเดลแบบแบ่งชั้น, กลยุทธ์ incremental, การทดสอบคุณภาพข้อมูล และคำถามสัมภาษณ์พร้อมตัวอย่างโค้ดสำหรับปี 2026