Pythonでconfluent-kafkaを使いこなす:毎秒数百万イベントを処理する秘訣

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

ブラックフライデーの物語とシステムダウンの悪夢

約半年前、あるECサイトのトラッキングシステムを構築するタスクを引き受けました。当初、構成は非常にシンプルでした。チームは「手堅い」方法を選択し、ユーザーが商品をクリックするたびに、PythonスクリプトがPostgreSQLに直接レコードを書き込んでいました。

すべては順調でしたが、ブラックフライデー当日に事態は一変しました。トラフィックは秒間5,000リクエストまで急増。データベースはToo many connectionsエラーを吐き出し始めました。レイテンシは50msから15秒まで跳ね上がり、最終的にシステムはフリーズ。数千件の注文ログデータが消失してしまいました。その時、大量のイベントを処理する際にデータを直接書き込む(同期処理)のは致命的なミスであると痛感しました。

なぜPythonスクリプトは息切れしてしまうのか?

システムを分析した結果、旧スクリプトが失敗した3つの核心的な問題が見えてきました:

  • データベースのボトルネック: PostgresやMySQLは、バッファなしで秒間10,000件のレコードを処理するようには設計されていません。
  • ドミノ倒し効果: データベースが1分間メンテナンスに入るだけで、データの送り先を失った前段のトラッキングスクリプトはすべて停止してしまいます。
  • 言語の壁: 純粋なPython製のkafka-pythonライブラリはインストールが簡単ですが、GIL(Global Interpreter Lock)の制限を受けるため、高負荷時の重いI/O処理が極端に遅くなります。

各オプションの比較:最適な選択肢はどれか?

現在主流となっている3つのアプローチを1週間かけてベンチマークしました。

1. Redisをキューとして使用する

Redisは非常に高速ですが、ストリーミングプラットフォームというよりは一時的なストレージに近いです。永続化(Persistence)メカニズムが不足しており、Consumerに障害が発生した際にデータを「読み直す」ことができません。

2. kafka-pythonライブラリ

pip installだけで済むため定番の選択肢ですが、秒間50,000メッセージでテストしたところ、CPU使用率は常に90-95%に達しました。過酷な本番環境で動かすには力不足です。

3. confluent-kafkaライブラリ

これが最終的な私の選択です。librdkafka(C言語製)をベースに構築されているため、ネットワーク処理をPythonのメインスレッドから分離でき、純粋なPython版に比べてスループットが3〜4倍向上します。

confluent-kafkaの実践的な実装

単にインストールしてデフォルト設定で動かすのではなく、最高のパフォーマンスを引き出すための私の設定方法を紹介します。

ステップ1:インストール

補足:Linux環境でこのライブラリをインストールする前に、librdkafkaをインストールしておく必要があります。

pip install confluent-kafka

ステップ2:本格的なProducerの作成

多くの人が犯す間違いは、メッセージを送信しっぱなしにすることです。データが確実にKafkaに届いたことを確認するために、delivery_reportをキャッチする必要があります。

from confluent_kafka import Producer
import json

conf = {
    'bootstrap.servers': "localhost:9092",
    'client.id': 'prod-collector-v1',
    'linger.ms': 20, # データを20ms間溜めて一括送信する
    'batch.size': 32768 # バッチサイズを32KBに増やす
}

producer = Producer(conf)

def on_delivery(err, msg):
    if err: print(f"エラーが発生しました: {err}")
    # 成功時はログを節約するために何もしない

def send_data(topic, payload):
    data = json.dumps(payload).encode('utf-8')
    producer.produce(topic, value=data, callback=on_delivery)
    producer.poll(0) # コールバックを即座にトリガーする

# 100万イベントの送信をシミュレート
for i in range(1000000):
    send_data("user_clicks", {"id": i, "type": "view"})

producer.flush() # バッファに残っているものを強制的に送信する

ヒント: スクリプトを終了する前にflush()関数を呼び出すことは必須です。これがないと、バッファに残っている最後の5〜10%のデータが失われる可能性があります。

ステップ3:堅牢なConsumer

Consumerに関しては安定性を優先します。決済のような重要なタスクを処理する場合は、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エラー: {msg.error()}")
            continue

        # ここでデータ処理ロジックを記述
        print(f"受信済み: {msg.value().decode('utf-8')}")
finally:
    consumer.close()

本番運用6ヶ月で得た3つの貴重な教訓

  1. JSONはパフォーマンスの敵: 100k msg/sに達すると、JSONのパースだけで膨大なCPUを消費します。私はProtobufに移行しました。これによりメッセージサイズを40%削減でき、帯域幅も大幅に節約できました。
  2. モニタリングを忘れない: 常にconsumer_lagを監視する必要があります。ラグが増え続ける場合は、パーティション数を増やし、Consumerのインスタンスを追加して負荷を分散させましょう。ノードを追加するだけで解決するなら、無理にコードの最適化に時間を費やす必要はありません。
  3. 能動的なエラー処理: Kafkaブローカーはいつでもダウンする可能性があります。produceコマンドをtry-exceptブロックで囲み、貴重なデータの損失を防ぐために、ローカルファイルへの一時書き出しなどのバックアッププランを用意しておきましょう。

confluent-kafkaへの移行により、追加のサーバー費用をかけずに大規模なセール期間を乗り切ることができました。本格的なリアルタイムアプリを構築するなら、今日から使い始めることをお勧めします。

Share: