Làm chủ confluent-kafka trong Python: Bí kíp xử lý hàng triệu Event mỗi giây

Python tutorial - IT technology blog
Python tutorial - IT technology blog

Câu chuyện Black Friday và cơn ác mộng sập hệ thống

Khoảng nửa năm trước, mình nhận task xây dựng hệ thống tracking cho một sàn TMĐT. Ban đầu, mọi thứ khá đơn giản. Team chọn cách “ăn chắc mặc bền”: Mỗi khi user click vào sản phẩm, script Python sẽ ghi thẳng một record vào PostgreSQL.

Mọi thứ vẫn êm đềm cho đến ngày Black Friday. Traffic tăng vọt lên mức 5.000 request/giây. Database bắt đầu gào thét với lỗi Too many connections. Latency nhảy vọt từ 50ms lên tận 15 giây. Cuối cùng, hệ thống treo cứng, khiến team mất trắng dữ liệu log của cả nghìn đơn hàng. Lúc đó mình mới thấm: Ghi dữ liệu trực tiếp (Synchronous) khi xử lý lượng event lớn là một sai lầm chí mạng.

Tại sao script Python của bạn thường bị hụt hơi?

Sau khi mổ xẻ hệ thống, mình nhận ra 3 vấn đề cốt lõi khiến script cũ thất bại:

  • Nghẽn cổ chai tại Database: Postgres hay MySQL không sinh ra để nuốt 10.000 bản ghi mỗi giây mà không có bộ đệm.
  • Hiệu ứng Domino: Chỉ cần Database bảo trì 1 phút, toàn bộ script tracking phía trước sẽ chết đứng vì không có chỗ đẩy data vào.
  • Rào cản ngôn ngữ: Thư viện kafka-python thuần Python rất dễ cài. Tuy nhiên, nó bị dính Global Interpreter Lock (GIL), khiến việc xử lý I/O nặng trở nên cực kỳ chậm chạp khi load cao.

Cân nhắc giữa các phương án: Đâu là lựa chọn tối ưu?

Mình đã dành 1 tuần để benchmark 3 hướng tiếp cận phổ biến nhất hiện nay.

1. Dùng Redis làm hàng chờ

Redis cực nhanh, nhưng nó giống như một kho chứa tạm thời hơn là một nền tảng streaming. Nó thiếu cơ chế lưu trữ lâu dài (Persistence) và không thể “đọc lại” dữ liệu nếu Consumer gặp sự cố.

2. Thư viện kafka-python

Đây là lựa chọn quốc dân vì chỉ cần pip install là xong. Nhưng khi mình test với 50.000 messages/giây, CPU luôn nhảy lên mức 90-95%. Nó không đủ độ lỳ để chạy trong môi trường production khắc nghiệt.

3. Thư viện confluent-kafka

Đây là lựa chọn sau cùng của mình. Nó được xây dựng trên librdkafka (viết bằng C). Nhờ đó, nó tách biệt việc xử lý network ra khỏi luồng chính của Python, giúp throughput tăng gấp 3-4 lần so với bản thuần Python.

Triển khai thực tế với confluent-kafka

Đừng chỉ cài đặt rồi chạy theo default. Dưới đây là cách mình cấu hình để đạt hiệu năng tốt nhất.

Bước 1: Cài đặt

Lưu ý nhỏ: Bạn cần cài đặt librdkafka trước khi install thư viện này trên môi trường Linux.

pip install confluent-kafka

Bước 2: Viết Producer “xịn”

Sai lầm của nhiều người là gửi tin nhắn xong… để đó. Bạn phải bắt được delivery_report để biết chắc chắn data đã nằm an toàn trong Kafka.

from confluent_kafka import Producer
import json

conf = {
    'bootstrap.servers': "localhost:9092",
    'client.id': 'prod-collector-v1',
    'linger.ms': 20, # Gom data trong 20ms để gửi một mẻ
    'batch.size': 32768 # Tăng kích thước batch lên 32KB
}

producer = Producer(conf)

def on_delivery(err, msg):
    if err: print(f"Lỗi rồi: {err}")
    # Thành công thì im lặng để tiết kiệm log

def send_data(topic, payload):
    data = json.dumps(payload).encode('utf-8')
    producer.produce(topic, value=data, callback=on_delivery)
    producer.poll(0) # Trigger callback ngay lập tức

# Giả lập đẩy 1 triệu event
for i in range(1000000):
    send_data("user_clicks", {"id": i, "type": "view"})

producer.flush() # Ép buộc gửi nốt những gì còn trong buffer

Mẹo nhỏ: Hàm flush() là bắt buộc trước khi tắt script. Thiếu nó, bạn sẽ mất khoảng 5-10% dữ liệu cuối cùng đang nằm trong bộ nhớ đệm.

Bước 3: Consumer bền bỉ

Với Consumer, mình ưu tiên sự ổn định. Nếu xử lý các task quan trọng như thanh toán, hãy tắt auto.commit.

from confluent_kafka import Consumer

conf = {
    'bootstrap.servers': "localhost:9092",
    'group.id': "analytics-group-1",
    'auto.offset.reset': 'earliest'
}

consumer = Consumer(conf)
consumer.subscribe(['user_clicks'])

try:
    while True:
        msg = consumer.poll(1.0)
        if msg is None: continue
        if msg.error():
            print(f"Consumer error: {msg.error()}")
            continue

        # Logic xử lý data ở đây
        print(f"Đã nhận: {msg.value().decode('utf-8')}")
finally:
    consumer.close()

3 bài học đắt giá sau 6 tháng lên Production

  1. JSON là kẻ thù của hiệu năng: Khi đạt ngưỡng 100k msg/s, việc parse JSON ngốn cực nhiều CPU. Mình đã chuyển sang Protobuf, giúp giảm 40% size tin nhắn và tiết kiệm đáng kể băng thông.
  2. Đừng quên Monitoring: Luôn phải theo dõi consumer_lag. Nếu lag tăng liên tục, hãy tăng số lượng Partition và chạy thêm instance Consumer để chia tải. Đừng cố tối ưu code khi chỉ cần thêm node là xong.
  3. Xử lý lỗi chủ động: Kafka broker có thể die bất cứ lúc nào. Hãy bọc lệnh produce trong khối try-except và có phương án dự phòng (như ghi tạm ra file local) để tránh mất dữ liệu quý giá.

Việc chuyển sang confluent-kafka giúp hệ thống của mình trụ vững qua các đợt sale lớn mà không tốn thêm chi phí server. Nếu bạn định làm app real-time nghiêm túc, hãy bắt đầu với nó ngay hôm nay.

Share: