---
title: コンシューマーグループと協調動作
url: https://doc.liz6.com/ja/distributed-systems/07-messages-and-streams/03-consumer-groups-and-collaboration
locale: ja
area: distributed-systems
tags:
- distributed-systems
- messages-and-streams
date: 2026-06-30
modified: 2026-07-19
description: コンシューマーグループは、パーティション所有権に対するコンシューマー間の合意形成であり、本質的には分散合意プロトコルの一種です。イグリーvs協調的なリバランスは、リバランス時にグループ全体が消費を停止するか継続するかを決定します。
---

# コンシューマーグループと協調動作

> コンシューマーグループは、パーティション所有権に対するコンシューマー間の合意形成であり、本質的には分散合意プロトコルの一種です。イグリーvs協調的なリバランスは、リバランス時にグループ全体が消費を停止するか継続するかを決定します。

## 概要

[Kafkaアーキテクチャ](/distributed-systems/07-messages-and-streams/02-kafka-architecture.md)では、パーティションとISR（In-Sync Replicas）について解説しました。これらはストレージ/レプリケーション層における分散メカニズムです。しかし、コンシューマー側にも分散協調の問題があります。複数のコンシューマーが1つのトピックのパーティションを共有する場合、どのパーティションをどのコンシューマーが消費するのか？ コンシューマーがダウンした場合、どのように割り当てを再分配するのか？ 再分配中は消費を停止するのか継続するのか？ これらの問題を解決するのが**コンシューマーグループ（consumer group）**です。本質的には「グループ内のコンシューマーによるパーティション所有権の合意形成」であり、紛れもない分散合意プロトコルの一種です。

## コンシューマーグループの動作原理

Kafkaの消費モデル：**1つのパーティションは、1つのコンシューマーグループ内では1つのコンシューマーのみが消費できます。** これが「パーティション内のメッセージ順序が保証される」ための前提条件です。もし2つのコンシューマーが同時に同じパーティションを読み取ると、それぞれが認識するオフセットの範囲が不連続になり、順序が乱れてしまいます。

<svg viewBox="0 0 720 340" xmlns="http://www.w3.org/2000/svg" font-family="-apple-system,'Source Han Sans CN','Microsoft YaHei',sans-serif" role="img" aria-label="コンシューマーグループの割り当て前後の比較: 6つのパーティションを3つのコンシューマーに割り当て">
  <defs>
    <marker id="cga1" markerWidth="10" markerHeight="8" refX="8" refY="3" orient="auto"><path d="M0,0 L8,3 L0,6 Z" fill="#475569"/></marker>
  </defs>
  <rect width="720" height="340" fill="#ffffff"/>
  <text x="360" y="28" text-anchor="middle" font-size="17" font-weight="700" fill="#1f2933">コンシューマーグループの割り当て: 6つのパーティションを3つのコンシューマーに分配</text>
  <line x1="360" y1="44" x2="360" y2="290" stroke="#e2e8f0" stroke-width="1"/>
  <text x="190" y="58" text-anchor="middle" font-size="13" font-weight="700" fill="#64748b">割り当て前（リバランス中）</text>
  <text x="530" y="58" text-anchor="middle" font-size="13" font-weight="700" fill="#4f46e5">割り当て後</text>

  <rect x="52" y="70" width="36" height="40" rx="4" fill="#e2e8f0"/>
  <text x="70" y="86" text-anchor="middle" font-size="11" font-weight="600" fill="#475569">P0</text>
  <text x="70" y="102" text-anchor="middle" font-size="10" fill="#94a3b8">??</text>
  <rect x="92" y="70" width="36" height="40" rx="4" fill="#e2e8f0"/>
  <text x="110" y="86" text-anchor="middle" font-size="11" font-weight="600" fill="#475569">P1</text>
  <text x="110" y="102" text-anchor="middle" font-size="10" fill="#94a3b8">??</text>
  <rect x="132" y="70" width="36" height="40" rx="4" fill="#e2e8f0"/>
  <text x="150" y="86" text-anchor="middle" font-size="11" font-weight="600" fill="#475569">P2</text>
  <text x="150" y="102" text-anchor="middle" font-size="10" fill="#94a3b8">??</text>
  <rect x="172" y="70" width="36" height="40" rx="4" fill="#e2e8f0"/>
  <text x="190" y="86" text-anchor="middle" font-size="11" font-weight="600" fill="#475569">P3</text>
  <text x="190" y="102" text-anchor="middle" font-size="10" fill="#94a3b8">??</text>
  <rect x="212" y="70" width="36" height="40" rx="4" fill="#e2e8f0"/>
  <text x="230" y="86" text-anchor="middle" font-size="11" font-weight="600" fill="#475569">P4</text>
  <text x="230" y="102" text-anchor="middle" font-size="10" fill="#94a3b8">??</text>
  <rect x="252" y="70" width="36" height="40" rx="4" fill="#e2e8f0"/>
  <text x="270" y="86" text-anchor="middle" font-size="11" font-weight="600" fill="#475569">P5</text>
  <text x="270" y="102" text-anchor="middle" font-size="10" fill="#94a3b8">??</text>
  <text x="190" y="130" text-anchor="middle" font-size="11" fill="#94a3b8">コンシューマーはまだどのパーティションも保持していない</text>

  <rect x="65" y="150" width="70" height="30" rx="6" fill="#f8fafc" stroke="#cbd5e1" stroke-dasharray="4 3"/>
  <text x="100" y="169" text-anchor="middle" font-size="12" fill="#94a3b8">C1</text>
  <rect x="155" y="150" width="70" height="30" rx="6" fill="#f8fafc" stroke="#cbd5e1" stroke-dasharray="4 3"/>
  <text x="190" y="169" text-anchor="middle" font-size="12" fill="#94a3b8">C2</text>
  <rect x="245" y="150" width="70" height="30" rx="6" fill="#f8fafc" stroke="#cbd5e1" stroke-dasharray="4 3"/>
  <text x="280" y="169" text-anchor="middle" font-size="12" fill="#94a3b8">C3</text>
  <text x="190" y="200" text-anchor="middle" font-size="11" fill="#94a3b8">グループコーディネーターが新しい割り当てを計算するのを待機中</text>

  <rect x="380" y="72" width="36" height="28" rx="4" fill="#0d9488"/>
  <text x="398" y="91" text-anchor="middle" font-size="10" fill="#ffffff">P0</text>
  <rect x="420" y="72" width="36" height="28" rx="4" fill="#0d9488"/>
  <text x="438" y="91" text-anchor="middle" font-size="10" fill="#ffffff">P1</text>
  <line x1="460" y1="86" x2="492" y2="86" stroke="#475569" stroke-width="1.6" marker-end="url(#cga1)"/>
  <rect x="496" y="72" width="150" height="28" rx="6" fill="#4f46e5"/>
  <text x="571" y="91" text-anchor="middle" font-size="13" font-weight="700" fill="#ffffff">C1</text>

  <rect x="380" y="120" width="36" height="28" rx="4" fill="#0d9488"/>
  <text x="398" y="139" text-anchor="middle" font-size="10" fill="#ffffff">P2</text>
  <rect x="420" y="120" width="36" height="28" rx="4" fill="#0d9488"/>
  <text x="438" y="139" text-anchor="middle" font-size="10" fill="#ffffff">P3</text>
  <line x1="460" y1="134" x2="492" y2="134" stroke="#475569" stroke-width="1.6" marker-end="url(#cga1)"/>
  <rect x="496" y="120" width="150" height="28" rx="6" fill="#4f46e5"/>
  <text x="571" y="139" text-anchor="middle" font-size="13" font-weight="700" fill="#ffffff">C2</text>

  <rect x="380" y="168" width="36" height="28" rx="4" fill="#0d9488"/>
  <text x="398" y="187" text-anchor="middle" font-size="10" fill="#ffffff">P4</text>
  <rect x="420" y="168" width="36" height="28" rx="4" fill="#0d9488"/>
  <text x="438" y="187" text-anchor="middle" font-size="10" fill="#ffffff">P5</text>
  <line x1="460" y1="182" x2="492" y2="182" stroke="#475569" stroke-width="1.6" marker-end="url(#cga1)"/>
  <rect x="496" y="168" width="150" height="28" rx="6" fill="#4f46e5"/>
  <text x="571" y="187" text-anchor="middle" font-size="13" font-weight="700" fill="#ffffff">C3</text>

  <rect x="60" y="300" width="600" height="30" rx="8" fill="#eef2ff" stroke="#c7d2fe"/>
  <text x="360" y="320" text-anchor="middle" font-size="12.5" fill="#3730a3">1つのパーティションは1つのグループ内で1つのコンシューマーのみが消費できる。6つのパーティションが並列処理の上限を決定し、7番目のコンシューマーが参加するとアイドル状態になる。</text>
</svg>

重要なのは、パーティション数が並列処理の上限を決定することです。パーティションが6つある場合、最大で6つのコンシューマーが利用可能です。7番目のコンシューマーが参加しても、割り当てられるパーティションがないため**アイドル状態**になります。**消費の並列度を拡張するには、コンシューマーを追加するのではなく、パーティションを追加する必要があります。**

## リバランス（再均衡）: 所有権の合意形成

コンシューマーがグループに参加または離脱すると、リバランスがトリガーされ、パーティションの所有権が再分配されます。これはコンシューマーグループのメカニズムにおいて、最も重要かつ問題が発生しやすい部分です。

### イグリーリバランス（旧プロトコル）: 停止して再分配

```
1. 新しいコンシューマー C4 が参加 → グループコーディネーターが全コンシューマーに通知: 現在の割り当てを解放
2. 全コンシューマーが消費を停止し、パーティションの所有権を解放
3. 再分配: C1←P0,P1, C2←P2, C3←P3,P4, C4←P5
4. 全コンシューマーが最後にコミットしたオフセットから消費を再開
```

このプロセスは **stop-the-world（ワールドストップ）** と呼ばれます。ステップ2-3の間、**コンシューマーグループ全体が消費を停止**し、レイテンシーが蓄積します。レイテンシーに敏感なストリーム処理アプリケーションにとって、リバランスは毎回の事故です。

### 協調的リバランス（新プロトコル、KIP-429）: 漸進的

```
1. C4 が参加 → コーディネーターが新しい割り当てを計算し、「パーティションを解放する必要がある」コンシューマーのみ通知
2. 影響を受けたコンシューマーは、奪われたパーティションを解放し、影響を受けないパーティションの消費を継続
3. 新しい所有者（C4）が解放されたパーティションを引き受け
4. この間、影響を受けないコンシューマーは停止しない
```

本質的な違い: イグリーでは全員が解放してから再分配を行うのに対し、協調的では影響を受けたもののみを変更し、他は停止しません。**「すべて停止」から「一部は停止しない」への変化**により、リバランスによるビジネスへの影響を大幅に軽減します。

### スティッキー割り当て: 既存の割り当てを維持する

リバランス後の割り当ては「均等性」だけでなく「最小変更」も考慮されます。もしあるコンシューマーが既に特定のパーティションを消費中で、解放する必要がない場合、その割り当ては**変更されません**。これを**スティッキー（sticky）**と呼び、不要な割り当て変更を避け、コンシューマーの状態遷移を減らします。

## オフセット管理: 消費進捗の「永続化ポインター」

コンシューマーがどこまで読んだかを記録する必要があります。これがオフセットです。管理方法には2つの種類があります。

| | 自動コミット | 手動コミット |
|---|---|---|
| 時期 | 時間間隔（例: 5秒ごと）で現在のオフセットを自動的にコミット | アプリケーションが `commitSync()` / `commitAsync()` を明示的に呼び出す |
| リスク | プロセス完了後、コミット間隔前にクラッシュ → 重複消費（at-least-once） | プロセス完了 → コミット → クラッシュ → 重複なし。ただし、プロセス完了後、コミット前にクラッシュ → 重複発生 |
| 適用 | 少量の重複を許容できる場合 | 細かな制御が必要な場合（例: exactly-once 消費フロー） |

**オフセットの保存場所**: Kafkaの内部トピック `__consumer_offsets`（マルチパーティション、レプリケーションあり）です。コンシューマーグループのオフセットコミットは、このトピックへの書き込みメッセージとして記録されます。つまり、Kafkaは自前のストレージメカニズムを使用して消費進捗を保存しています。

## exactly-once 消費: 消費側トランザクションの結合

[メッセージセマンティクス](/distributed-systems/07-messages-and-streams/01-message-semantics.md)では、プロデューサー側の exactly-once（冪等プロデューサー + トランザクション）について解説しました。消費側で exactly-once を実現するには、**オフセットのコミット**と**処理結果の出力**を同じトランザクションに含めることが核心です。

```
コンシューマートランザクション:
  1. topic-A からメッセージを読み取る (offset=N)
  2. 処理 → 結果を topic-B に書き込み
  3. __consumer_offsets へ offset(N+1) をコミット
  ステップ2と3は同じKafkaトランザクション内で行われる:
    - トランザクション成功: 結果が topic-B に書き込まれ、offset が N+1 に進む → 消費が「完了した」とみなされる
    - トランザクション失敗/コンシューマークラッシュ: 結果は書き込まれず、offset は N にロールバック → メッセージは再消費される
```

重要な要件: 消費時の `isolation.level` を `read_committed` に設定する必要があります。これにより、コミット済みトランザクションのメッセージのみを読み取り、アボートされたトランザクションや未完成のdanglingトランザクションを無視します。

この一連の組み合わせ（冪等プロデューサー + トランザクションプロデューサー + `read_committed` コンシューマー + オフセットと出力が同じトランザクション内）により、Kafkaの exactly-once セマンティクスが構成されます。これは理論上の「完璧な僅か1回の処理」ではなく、工学的に「僅か1回の処理と同等の効果を持ち、重複も欠落も発生しない」ことを意味します。

## RabbitMQとの比較: キューモデル vs ログモデル

同じ章で解説されている別のアプローチ: RabbitMQはログではなく**キュー**を使用します。メッセージは確認されると**削除**されます（不変ログへの追加ではないため）、リプレイ/巻き戻しはできず、パーティション内のグローバル順序も保証されません。コンシューマーグループの協調も異なります。RabbitMQのコンシューマーはパーティション所有権を合意形成するのではなく、ブローカーがラウンドロビンまたは一貫性ハッシュ交換によってメッセージを配信します。これはよりシンプルですが、Kafkaのような「同じキーは同じパーティションへ」という順序保証はありません。

Kafkaのログモデルが選ばれる理由は、分散ストリーム処理において**リプレイ（replay）**と**グローバル順序**が必要だからです。これらはキューモデルでは実現できません。Kafkaのパーティションは、シャーディングされた不変でリプレイ可能なログそのものであり、コンシューマーグループはこのログの上で「どのパーティションを誰が読むか」を合意形成する層です。

## よくある問題と解決策

- **リバランスストーム**: コンシューマーの頻繁な参加/離脱（例: ハートビートタイムアウト閾値が短すぎる、処理時間が長すぎる → ダウンとみなされる）。毎回のリバランスでグループ全体が停止する。→ `max.poll.interval.ms` を十分に大きく設定する（最長処理時間より長く）、協調的リバランスを使用し、リバランス頻度を監視する。
- **パーティション数 < コンシューマー数**: 余分なコンシューマーがアイドル状態になる → パーティションを追加するか、アイドル状態を受け入れる（コンシューマーをホットスタンバイとして使用）。
- **消費遅延（lag）**: コンシューマーの処理速度がプロデューサーの書き込み速度に追いつかない → lagを監視（コンシューマーオフセット vs パーティションエンドオフセット）、コンシューマーまたはパーティションを追加する。
- **自動コミット + 非同期処理**: コンシューマーがメッセージをスレッドプールに投げてオフセットをコミット → 非同期スレッドの処理失敗時にメッセージは既に「消費済み」となる → 消失。→ 同期処理を行うか、非同期処理完了後に手動コミットを行う。

## 参考

- **Kafkaドキュメント**: Consumer Group Protocol、KIP-429 (Incremental Cooperative Rebalancing)

*Keywords: consumer group, rebalance, eager, cooperative (incremental), sticky assignment, stop-the-world, group coordinator, offset management, __consumer_offsets, exactly-once consumption, isolation.level, read_committed, transactional consumer, 消費遅延, consumer lag, rebalance storm, Kafka vs RabbitMQ, log vs queue*
