•15 min read

ApacheKafkaによるリアルタイムデータストリーミング:実践ガイド

ApacheKafkaによるリアルタイムデータストリーミング:実践ガイド

データがミリ秒単位で価値を失う世界では、バッチ処理だけではもはや十分ではありません。リアルタイムのデータストリーミングは、現代のデータアーキテクチャのバックボーンとなっており、Apache Kafkaは分散イベントストリーミングの議論の余地のない業界標準です。LinkedIn(誕生の地)、Netflix、Uber、Airbnbなどの企業で、1日に数兆のイベントを処理しています。

このガイドでは、本番環境のKafkaシステムを構築、運用、スケーリングするために必要なすべてを網羅しています。アーキテクチャ、Pythonプロデューサー/コンシューマー、パーティショニング戦略、Exactly-Onceセマンティクス、バックプレッシャー処理、ゼロダウンタイムスケーリングなどです。

Audio Briefing
0:00 / 0:00

Kafkaのアーキテクチャ

Kafkaは分散コミットログとして機能します。プロデューサーはイベント(メッセージ)を名前付きのトピックにパブリッシュします。コンシューマーはそれらのトピックを購読し、独自のペースでイベントを処理します。従来のメッセージキュー(RabbitMQ、SQS)が消費後にメッセージを削除するのとは異なり、Kafkaは設定された期間イベントを保持します。これにより、複数のコンシューマーグループが同じデータを独立して読み取ることができ、任意のオフセットからリプレイすることが可能になります。

主要なアーキテクチャコンポーネント:

コンポーネント役割
Broker単一のKafkaサーバー。トピックパーティションをディスクに保存します。
Topic名前付きのイベントストリーム。パーティションに分割されます。
Partition並列処理の単位。順序付けされた不変のログ。
Producerイベントをトピックに書き込みます。
Consumerイベントをトピックから読み取ります。
Consumer Groupトピックを協調して処理するコンシューマーのセット(各パーティションはグループ内の1つのメンバーによってのみ消費されます)。
ZooKeeper / KRaftメタデータ調整。KRaft(Kafka Raft)は、Kafka 3.3以降でZooKeeperに代わり、運用を簡素化します。
Advertisement

トピック、パーティション、順序付け

パーティショニングはKafkaの水平スケーリングメカニズムです。1つのトピックは1から数千のパーティションを持つことができます。各パーティションは、特定のブローカーに保存される、順序付けされた追記専用のログです。

**順序付けはパーティション内では保証されますが、パーティション間では保証されません。**これは重要な設計上の制約です。特定のユーザーのすべてのイベントを順番に処理する必要がある場合、一貫したパーティションキーを使用して、そのユーザーのすべてのイベントを同じパーティションにルーティングする必要があります。

# 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()
)

Kafkaのデフォルトのパーティショナーは、ターゲットパーティションを決定するためにmurmur2(key) % num_partitionsを使用します。

パーティションの数は? 一般的な経験則として、パーティションあたり約10 MB/秒のスループットを目指します。合計1 GB/秒を処理するトピックの場合、約100のパーティションが必要になります。パーティションが多いほど並列処理は増えますが、コーディネーターのオーバーヘッド(ファイルハンドルの増加、障害発生時のリーダー選出時間の延長)も大きくなります。

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

主要なプロデューサー設定の説明:

  • acks=all — 最も強力な耐久性保証。リーダーは、すべての同期レプリカ(ISR)が応答を承認するまで待機してから、プロデューサーに応答します。
  • enable.idempotence=True — プロデューサーのリトライ時の重複メッセージを防ぎます(acks=allが必要です)。
  • linger.ms + batch.size — 複数のレコードを単一のリクエストにバッチ処理することで、わずかなレイテンシーと引き換えに大幅に高いスループットを実現します。
  • compression.type=lz4 — LZ4は、CPU負荷が低く、JSONペイロードの圧縮率が高い、オールラウンドで最適なKafka圧縮方式です。

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**は、アットリーストワンス(またはExactly-Once)処理にとって非常に重要です。自動コミットが有効になっている場合、Kafkaは処理が成功したかどうかに関わらず、定期的にオフセットをコミットします。自動コミットと処理完了の間にクラッシュが発生すると、イベントが失われます。

Advertisement

コンシューマーのバックプレッシャー処理

バックプレッシャーは、コンシューマーがプロデューサーがイベントを書き込むよりも遅くイベントを処理するときに発生します。これを処理しないと、無制限のコンシューマーラグ、メモリプレッシャー、そして最終的にはOOMクラッシュが発生します。

バックプレッシャーを処理するためのパターン:

1. 低速コンシューマー:並列処理を増やす

パーティションとコンシューマーを増やします(グループあたり最大num_partitionsコンシューマー):

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

次に、コンシューマーグループをスケーリングします。各コンシューマーは24 / num_consumersパーティションを処理します。

2. 制限付きキューによる非同期処理

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. デッドレターキュー(DLQ)

繰り返し処理に失敗するイベントは、メインのトピックをブロックするのではなく、DLQにルーティングします。

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"})

ログ圧縮

Kafkaは2つの保持ポリシーをサポートしています。

  • 時間/サイズベースの保持 — N日より古いセグメントまたはNバイトより大きいセグメントを削除します。デフォルトの動作です。
  • ログ圧縮 — キーごとに最新のイベントのみを保持します。チェンジログトピックやマテリアライズドビューに最適です。

トピックでログ圧縮を有効にするには:

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"

圧縮されたuser-profilesトピックには、常に各user_idキーの最新のプロファイルが含まれています。これは、キャッシュのブートストラップや読み取りモデルの再構築に最適です。

Exactly-Onceセマンティクス(EOS)

Kafka 3.0以降では、Kafka Streamsを使用するか、トランザクションプロデューサーとコンシューマーを手動で調整することで、エンドツーエンドのExactly-Once処理をサポートしています。

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は、アットリーストワンスと比較してスループットのオーバーヘッド(約20〜30%)があります。すべてのイベントストリームではなく、重複処理が正確性の問題(金融取引、在庫更新など)を引き起こす場合にのみ使用してください。

ゼロダウンタイムクラスタースケーリング

ダウンタイムなしでKafkaクラスターをスケーリングするには、次の手順が必要です。

ブローカーの追加:

# 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

スロットル(--throttle)は非常に重要です。これがないと、パーティションの再割り当てがネットワークを飽和させ、プロデューサー/コンシューマーを停止させてしまう可能性があります。

ローリングリスタート(設定変更またはアップグレードの場合):

  1. 開始前にmin.insync.replicas=2とreplication.factor=3を設定します。
  2. ブローカーを1つずつ再起動します。
  3. 再起動中のブローカーからパーティションリーダーが再割り当てされるのを待ってから、次に進みます。
  4. ISRサイズを監視します。完全なレプリケーションファクターに戻ってからのみ続行します。

監視:監視すべき主要なメトリクス

メトリクスツールアラートしきい値
kafka.consumer.lagJMX / Burrow継続的に1万メッセージ以上
kafka.broker.underReplicatedPartitionsJMX0より大きい(ブローカー障害またはネットワークの問題を示唆)
kafka.broker.requestHandlerAvgIdlePercentJMX30%未満(ブローカーCPU飽和)
kafka.producer.record-error-rateJMX0より大きい(プロデューサー障害)
ブローカーごとのディスク使用量Prometheus/node_exporter80%以上

クラスターのWebダッシュボードには、Kafka UIまたはAKHQを使用してください。

Kafkaと代替手段の使い分け

ユースケース最適な選択肢理由
高スループットイベントログKafka水平スケーリング、保持、リプレイ
低レイテンシーのタスクキューRabbitMQ / Redis Streams低いオーバーヘッド、簡単な運用
シンプルなpub/subRedis Pub/Sub永続性不要
クラウドネイティブなサーバーレスAWS Kinesis / Google Pub/Subマネージド、運用負担なし
ML特徴量パイプラインKafka + Flink大規模なステートフルストリーム処理

よくある質問

適切なパーティション数はどのように選択しますか? max(target_throughput_MB_s / 10, num_consumers)から始めましょう。パーティションは後でいつでも増やすことができます(kafka-topics.sh --alterを実行)。しかし、トピックを再作成しない限り減らすことはできません。パーティションが多すぎると、ファイルハンドルの数が増え、障害発生時のリーダー選出時間が長くなります。

コンシューマーグループとコンシューマーの違いは何ですか? コンシューマーグループは、トピックに対する論理的なサブスクライバーです。Kafkaは、グループ内の各パーティションが1つのコンシューマーによってのみ消費されることを保証します。複数のグループが同じトピックを独立して消費できます。これは、1つのイベントストリームから分析、MLパイプライン、監査ログにファンアウトする方法です。

Kafka Connectを使用すべきか、独自のコンシューマーを記述すべきか? 標準的なソース/シンクパターン(データベース → Kafka、Kafka → S3、Kafka → Elasticsearch)には、既存のコネクターを備えたKafka Connectを使用してください。独自のコンシューマーは、カスタムビジネスロジックの場合にのみ記述してください。

イベントはどのくらいの期間保持すべきですか? 監査/コンプライアンスの場合:90〜365日。リアルタイム処理パイプラインの場合:通常7日で十分です。圧縮されたトピック(チェンジログ/状態)の場合:無期限(圧縮によりサイズを管理可能に保ちます)。

まとめ

Apache Kafkaは、データ生成とデータ消費を大規模に分離できるため、非常に強力です。順序付けのためのスマートなパーティショニング、耐久性のための冪等なプロデューサー、正確性のための手動オフセットコミット、状態のためのログ圧縮といった主要なパターンが、本番環境のKafkaデプロイメントをおもちゃのような例と区別するものです。

シンプルに始めましょう。1つのトピック、1つのプロデューサー、1つのコンシューマーグループから。スループットの限界に達したらパーティションを追加します。同じデータを複数のシステムで必要とする場合はコンシューマーグループを追加します。ビジネスロジックが要求する場合にのみEOSを導入します。

こちらもおすすめです

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
Kafkaを用いたイベント駆動型アーキテクチャ
tech

Kafkaを用いたイベント駆動型アーキテクチャ

ApacheKafkaを使用して回復力のあるイベント駆動型システムを構築し、パーティションキーイング、コンシューマーグループのリバランス、およびexactly-once処理セマンティクスを習得しましょう。

Read more