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

Table of Contents
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ỉ.
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ần | Vai trò |
|---|---|
| Broker | Một máy chủ Kafka duy nhất. Lưu trữ các phân vùng topic trên đĩa. |
| Topic | Mộ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. |
| Producer | Ghi sự kiện vào một topic. |
| Consumer | Đọc sự kiện từ một topic. |
| Consumer Group | Mộ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. |
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ầuacks=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.
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):
- Đặt
min.insync.replicas=2vàreplication.factor=3trước khi bắt đầu - Khởi động lại từng broker một
- Đợ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
- 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.lag | JMX / Burrow | > 10k thông điệp trong thời gian dài |
kafka.broker.underReplicatedPartitions | JMX | > 0 (cho biết lỗi broker hoặc sự cố mạng) |
kafka.broker.requestHandlerAvgIdlePercent | JMX | < 30% (bão hòa CPU của broker) |
kafka.producer.record-error-rate | JMX | > 0 (lỗi producer) |
| Sử dụng đĩa trên mỗi broker | Prometheus/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ụng | Lựa chọn tốt nhất | Lý do |
|---|---|---|
| Log sự kiện thông lượng cao | Kafka | Khả năng mở rộng ngang, lưu giữ, phát lại |
| Hàng đợi tác vụ độ trễ thấp | RabbitMQ / Redis Streams | Chi phí thấp hơn, vận hành đơn giản hơn |
| Pub/sub đơn giản | Redis Pub/Sub | Không cần lưu trữ |
| Serverless gốc đám mây | AWS Kinesis / Google Pub/Sub | Được quản lý, không gánh nặng vận hành |
| Pipelines tính năng ML | Kafka + Flink | Xử 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
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

BigQuery + Cloud Run: Xây Dựng Pipeline Nhập Dữ Liệu Serverless Cho Production
Cẩm nang cấp production về nhập dữ liệu serverless trên Google Cloud: BigQuery Storage Write API, chiến lược phân vùng và phân cụm, bộ nhận FastAPI async trên Cloud Run, Terraform đầy đủ, phân tích chi phí thực tế, và những chế độ lỗi gọi bạn lúc 3 giờ sáng.
Read moreClickHouse vs DuckDB: So sánh sâu về kiến trúc và benchmark trong môi trường production
So sánh hai engine cơ sở dữ liệu phân tích dạng cột hàng đầu. Khám phá khi nào nên sử dụng vectorization nhúng với DuckDB và khi nào dùng OLAP phân tán thời gian thực với ClickHouse.
Read moreFastAPI vs Litestar: So sánh Benchmark Production và Kiến trúc Microservice Thông lượng Cao
Một bài so sánh khách quan, dựa trên benchmark giữa FastAPI và Litestar. Khám phá hiệu năng ASGI, kiến trúc dependency injection, tốc độ serialization, và typing cho OpenAPI.
Read more