•12 min read

Truyền dữ liệu thời gian thực với Apache Kafka: Hướng dẫn sản xuất

Truyền dữ liệu thời gian thực với Apache Kafka: Hướng dẫn sản xuất

Trong một thế giới mà dữ liệu mất giá trị từng mili giây, xử lý hàng loạt không còn đủ nữa. Truyền dữ liệu thời gian thực đã trở thành xương sống của kiến trúc dữ liệu hiện đại, và Apache Kafka là tiêu chuẩn công nghiệp không thể tranh cãi cho việc truyền sự kiện phân tán — xử lý hàng nghìn tỷ sự kiện mỗi ngày tại các công ty như LinkedIn (nơi nó ra đời), Netflix, Uber và Airbnb.

Hướng dẫn này bao gồm mọi thứ bạn cần để xây dựng, vận hành và mở rộng một hệ thống Kafka sản xuất: kiến trúc, producer/consumer Python, chiến lược phân vùng, ngữ nghĩa chính xác một lần, xử lý áp lực ngược và mở rộng không ngừng nghỉ.

Audio Briefing
0:00 / 0:00

Kiến trúc Kafka

Kafka hoạt động như một commit log phân tán. Producer xuất bản các sự kiện (thông điệp) đến các topic được đặt tên. Consumer đăng ký các topic đó và xử lý các sự kiện theo tốc độ của riêng chúng. Không giống như các hàng đợi thông điệp truyền thống (RabbitMQ, SQS) xóa thông điệp sau khi tiêu thụ, Kafka giữ lại các sự kiện trong một khoảng thời gian được cấu hình — cho phép nhiều nhóm consumer đọc cùng một dữ liệu một cách độc lập, và cho phép phát lại từ bất kỳ offset nào.

Các thành phần kiến trúc cốt lõi:

Thành phầnVai trò
BrokerMột máy chủ Kafka duy nhất. Lưu trữ các phân vùng topic trên đĩa.
TopicMột luồng sự kiện được đặt tên. Được chia thành các phân vùng.
PartitionĐơn vị song song. Một log có thứ tự, bất biến.
ProducerGhi sự kiện vào một topic.
ConsumerĐọc sự kiện từ một topic.
Consumer GroupMột tập hợp các consumer cùng nhau xử lý một topic (mỗi phân vùng được tiêu thụ bởi chính xác một thành viên).
ZooKeeper / KRaftĐiều phối siêu dữ liệu. KRaft (Kafka Raft) thay thế ZooKeeper trong Kafka 3.3+ để đơn giản hóa hoạt động.
Advertisement

Topics, Partitions và Thứ tự

Phân vùng là cơ chế mở rộng theo chiều ngang của Kafka. Một topic có thể có từ 1 đến hàng nghìn phân vùng. Mỗi phân vùng là một log có thứ tự, chỉ thêm vào, được lưu trữ trên một broker cụ thể.

Thứ tự được đảm bảo trong một phân vùng, không phải giữa các phân vùng. Đây là một ràng buộc thiết kế quan trọng. Nếu bạn cần tất cả các sự kiện cho một người dùng cụ thể được xử lý theo thứ tự, bạn phải định tuyến tất cả các sự kiện của người dùng đó đến cùng một phân vùng bằng cách sử dụng một khóa phân vùng nhất quán:

# Producer: route by user_id to guarantee per-user ordering
producer.produce(
    topic="user-events",
    key=str(user_id).encode(),   # same key → same partition
    value=json.dumps(event).encode()
)

Bộ phân vùng mặc định của Kafka sử dụng murmur2(key) % num_partitions để xác định phân vùng đích.

Bao nhiêu phân vùng? Một quy tắc chung: nhắm mục tiêu thông lượng ~10 MB/s cho mỗi phân vùng. Đối với một topic xử lý tổng cộng 1 GB/s, bạn sẽ muốn ~100 phân vùng. Nhiều phân vùng hơn = song song hơn nhưng chi phí điều phối cao hơn (nhiều file handle hơn, thời gian bầu cử leader lâu hơn khi xảy ra lỗi).

Xây dựng Producer Python

from confluent_kafka import Producer
import json
import time

conf = {
    "bootstrap.servers": "kafka1:9092,kafka2:9092,kafka3:9092",
    "acks": "all",                    # wait for all in-sync replicas to acknowledge
    "retries": 5,
    "retry.backoff.ms": 200,
    "compression.type": "lz4",        # compress batches — massive throughput boost
    "linger.ms": 5,                    # batch for up to 5ms before sending
    "batch.size": 65536,               # 64 KB batch size
    "enable.idempotence": True,        # exactly-once producer semantics
}

producer = Producer(conf)

def delivery_report(err, msg):
    if err is not None:
        print(f"Delivery failed for {msg.key()}: {err}")
    else:
        print(f"Delivered to {msg.topic()} [{msg.partition()}] @ offset {msg.offset()}")

# Produce a batch of events
for i in range(1000):
    event = {"user_id": i % 100, "action": "page_view", "ts": time.time()}
    producer.produce(
        topic="user-events",
        key=str(event["user_id"]),
        value=json.dumps(event),
        callback=delivery_report
    )
    producer.poll(0)    # trigger callbacks for previously sent messages

producer.flush()        # wait for all in-flight messages to be delivered

Giải thích các cài đặt producer chính:

  • acks=all — Đảm bảo độ bền mạnh nhất. Leader đợi tất cả các bản sao đồng bộ (ISR) xác nhận trước khi phản hồi producer.
  • enable.idempotence=True — Ngăn chặn các thông điệp trùng lặp khi producer thử lại (yêu cầu acks=all).
  • linger.ms + batch.size — Đánh đổi một độ trễ nhỏ để có thông lượng cao hơn đáng kể bằng cách nhóm nhiều bản ghi vào một yêu cầu duy nhất.
  • compression.type=lz4 — LZ4 là nén Kafka tốt nhất toàn diện: nhanh về CPU và tỷ lệ cao cho các payload JSON.

Xây dựng Consumer Python

from confluent_kafka import Consumer, KafkaException
import json

conf = {
    "bootstrap.servers": "kafka1:9092,kafka2:9092,kafka3:9092",
    "group.id": "user-event-processor-v1",
    "auto.offset.reset": "earliest",     # start from the beginning if no committed offset
    "enable.auto.commit": False,          # manual commit for exactly-once processing
    "max.poll.interval.ms": 300000,       # 5 min max processing time per batch
    "session.timeout.ms": 30000,
}

consumer = Consumer(conf)
consumer.subscribe(["user-events"])

try:
    while True:
        msg = consumer.poll(timeout=1.0)

        if msg is None:
            continue
        if msg.error():
            raise KafkaException(msg.error())

        event = json.loads(msg.value().decode("utf-8"))

        # Process the event
        process_event(event)

        # Commit offset only AFTER successful processing
        consumer.commit(asynchronous=False)

except KeyboardInterrupt:
    pass
finally:
    consumer.close()

enable.auto.commit=False rất quan trọng đối với xử lý ít nhất một lần (hoặc chính xác một lần). Với tính năng tự động commit được bật, Kafka commit offset định kỳ bất kể quá trình xử lý của bạn có thành công hay không — một sự cố giữa tự động commit và hoàn thành xử lý có nghĩa là bạn mất sự kiện.

Advertisement

Xử lý áp lực ngược của Consumer

Áp lực ngược xảy ra khi consumer của bạn xử lý các sự kiện chậm hơn producer ghi chúng. Nếu không xử lý điều này, bạn sẽ gặp tình trạng lag consumer không giới hạn, áp lực bộ nhớ và cuối cùng là các sự cố OOM.

Các mẫu để xử lý áp lực ngược:

1. Consumer chậm: tăng tính song song

Thêm nhiều phân vùng và nhiều consumer hơn (tối đa num_partitions consumer cho mỗi nhóm):

kafka-topics.sh --alter --topic user-events --partitions 24 \
  --bootstrap-server kafka1:9092

Sau đó mở rộng nhóm consumer của bạn: mỗi consumer xử lý 24 / num_consumers phân vùng.

2. Xử lý bất đồng bộ với hàng đợi có giới hạn

import asyncio
from asyncio import Queue

async def consume(consumer: Consumer, queue: Queue):
    while True:
        msg = consumer.poll(timeout=0.1)
        if msg and not msg.error():
            await queue.put(msg)       # blocks if queue is full — natural backpressure
            
async def process(queue: Queue):
    while True:
        msg = await queue.get()
        event = json.loads(msg.value())
        await process_event_async(event)
        queue.task_done()

# Bounded queue: max 500 in-flight events
queue = Queue(maxsize=500)
await asyncio.gather(consume(consumer, queue), process(queue))

3. Hàng đợi thư chết (DLQ)

Đối với các sự kiện liên tục không xử lý được, hãy định tuyến chúng đến DLQ thay vì chặn topic chính:

def process_with_dlq(msg):
    for attempt in range(3):
        try:
            process_event(json.loads(msg.value()))
            return
        except Exception as e:
            if attempt == 2:
                # Send to DLQ with original headers + error metadata
                producer.produce("user-events.DLQ", value=msg.value(),
                                 headers={"error": str(e), "original-topic": "user-events"})

Nén Log

Kafka hỗ trợ hai chính sách lưu giữ:

  • Lưu giữ dựa trên thời gian/kích thước — Xóa các phân đoạn cũ hơn N ngày hoặc lớn hơn N byte. Hành vi mặc định.
  • Nén log — Chỉ giữ lại sự kiện gần đây nhất cho mỗi khóa. Lý tưởng cho các topic changelog và các chế độ xem được vật chất hóa.

Bật nén log trên một topic:

kafka-configs.sh --bootstrap-server kafka1:9092 \
  --entity-type topics \
  --entity-name user-profiles \
  --alter \
  --add-config "cleanup.policy=compact,min.cleanable.dirty.ratio=0.1,delete.retention.ms=86400000"

Một topic user-profiles được nén luôn chứa hồ sơ mới nhất cho mỗi khóa user_id — hoàn hảo để khởi tạo bộ nhớ cache hoặc xây dựng lại mô hình đọc.

Ngữ nghĩa chính xác một lần (EOS)

Kafka 3.0+ hỗ trợ xử lý chính xác một lần từ đầu đến cuối khi sử dụng Kafka Streams hoặc phối hợp thủ công các producer và consumer giao dịch:

from confluent_kafka import Producer

conf = {
    "bootstrap.servers": "kafka1:9092",
    "transactional.id": "my-producer-instance-1",   # unique per producer instance
    "enable.idempotence": True,
}

producer = Producer(conf)
producer.init_transactions()

try:
    producer.begin_transaction()
    
    # Read from consumer, process, write result — all in one transaction
    for msg in batch:
        result = transform(json.loads(msg.value()))
        producer.produce("processed-events", value=json.dumps(result))
    
    # Commit consumer offsets and producer messages atomically
    producer.send_offsets_to_transaction(consumer_offsets, consumer.consumer_group_metadata())
    producer.commit_transaction()
    
except Exception:
    producer.abort_transaction()
    raise

EOS có chi phí thông lượng cao hơn (~20-30%) so với ít nhất một lần. Sử dụng nó khi xử lý trùng lặp gây ra các vấn đề về tính đúng đắn (giao dịch tài chính, cập nhật hàng tồn kho) thay vì cho tất cả các luồng sự kiện.

Mở rộng cụm không ngừng nghỉ

Mở rộng cụm Kafka mà không ngừng nghỉ bao gồm:

Thêm broker:

# 1. Add new broker to cluster — it joins automatically
# 2. Reassign partitions to include new broker
kafka-reassign-partitions.sh \
  --bootstrap-server kafka1:9092 \
  --topics-to-move-json-file topics.json \
  --broker-list "1,2,3,4" \          # include new broker 4
  --generate

# 3. Execute the reassignment plan
kafka-reassign-partitions.sh \
  --bootstrap-server kafka1:9092 \
  --reassignment-json-file reassignment.json \
  --execute \
  --throttle 50000000                 # 50 MB/s replication throttle — don't saturate network

Việc điều tiết (--throttle) là rất quan trọng. Nếu không có nó, việc gán lại phân vùng có thể làm bão hòa mạng của bạn và làm đói các producer/consumer.

Khởi động lại luân phiên (đối với thay đổi cấu hình hoặc nâng cấp):

  1. Đặt min.insync.replicas=2 và replication.factor=3 trước khi bắt đầu
  2. Khởi động lại từng broker một
  3. Đợi các leader phân vùng gán lại khỏi broker đang khởi động lại trước khi tiếp tục
  4. Giám sát kích thước ISR — chỉ tiếp tục khi nó trở lại hệ số nhân bản đầy đủ

Giám sát: Các chỉ số chính cần theo dõi

Chỉ sốCông cụNgưỡng cảnh báo
kafka.consumer.lagJMX / Burrow> 10k thông điệp trong thời gian dài
kafka.broker.underReplicatedPartitionsJMX> 0 (cho biết lỗi broker hoặc sự cố mạng)
kafka.broker.requestHandlerAvgIdlePercentJMX< 30% (bão hòa CPU của broker)
kafka.producer.record-error-rateJMX> 0 (lỗi producer)
Sử dụng đĩa trên mỗi brokerPrometheus/node_exporter> 80%

Sử dụng Kafka UI hoặc AKHQ để có bảng điều khiển web cho cụm của bạn.

Khi nào nên sử dụng Kafka so với các lựa chọn thay thế

Trường hợp sử dụngLựa chọn tốt nhấtLý do
Log sự kiện thông lượng caoKafkaKhả năng mở rộng ngang, lưu giữ, phát lại
Hàng đợi tác vụ độ trễ thấpRabbitMQ / Redis StreamsChi phí thấp hơn, vận hành đơn giản hơn
Pub/sub đơn giảnRedis Pub/SubKhông cần lưu trữ
Serverless gốc đám mâyAWS Kinesis / Google Pub/SubĐược quản lý, không gánh nặng vận hành
Pipelines tính năng MLKafka + FlinkXử lý luồng có trạng thái ở quy mô lớn

Câu hỏi thường gặp

Làm cách nào để chọn đúng số lượng phân vùng? Bắt đầu với max(target_throughput_MB_s / 10, num_consumers). Bạn luôn có thể tăng phân vùng sau này (chạy kafka-topics.sh --alter), nhưng bạn không thể giảm chúng mà không tạo lại topic. Việc phân vùng quá mức làm tăng số lượng file handle và thời gian bầu cử leader khi xảy ra lỗi.

Sự khác biệt giữa nhóm consumer và consumer là gì? Một nhóm consumer là một subscriber logic cho một topic. Kafka đảm bảo mỗi phân vùng được tiêu thụ bởi chính xác một consumer trong một nhóm. Nhiều nhóm có thể tiêu thụ cùng một topic một cách độc lập — đây là cách bạn phân phối đến phân tích, pipelines ML và ghi nhật ký kiểm toán từ một luồng sự kiện.

Tôi nên sử dụng Kafka Connect hay tự viết consumer của riêng mình? Đối với các mẫu nguồn/đích tiêu chuẩn (cơ sở dữ liệu → Kafka, Kafka → S3, Kafka → Elasticsearch), hãy sử dụng Kafka Connect với một connector hiện có. Chỉ viết consumer của riêng bạn cho logic nghiệp vụ tùy chỉnh.

Tôi nên giữ lại các sự kiện trong bao lâu? Đối với kiểm toán/tuân thủ: 90–365 ngày. Đối với pipelines xử lý thời gian thực: 7 ngày thường là đủ. Đối với các topic được nén (changelog/state): vô thời hạn (nén giúp giữ kích thước có thể quản lý được).

Tổng kết

Apache Kafka mạnh mẽ chính xác vì nó tách rời việc sản xuất dữ liệu khỏi việc tiêu thụ dữ liệu ở quy mô lớn. Các mẫu chính — phân vùng thông minh để sắp xếp thứ tự, producer bất biến để đảm bảo độ bền, commit offset thủ công để đảm bảo tính đúng đắn và nén log cho trạng thái — là những gì phân biệt các triển khai Kafka sản xuất với các ví dụ thử nghiệm.

Bắt đầu đơn giản: một topic, một producer, một nhóm consumer. Thêm phân vùng khi bạn đạt giới hạn thông lượng. Thêm nhóm consumer khi bạn cần cùng một dữ liệu trong nhiều hệ thống. Chỉ áp dụng EOS khi logic nghiệp vụ của bạn yêu cầu.

Bạn cũng có thể thích

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement