こんにちは。GxPの佐野です。
グロースエクスパートナーズグループのリレーブログ企画9日目です。
前回の記事は「テスト自動化したいけど、何から始める?」を手探りで進めた記録でした。
まだご覧になっていない方は、ぜひそちらもチェックしてみてください!
はじめに
Apache Kafkaとは、「リアルタイムで大量のデータを処理する分散メッセージングシステム」のことです。
今回のシステムでは複数の外部APIを呼び出して処理する必要がありました。
同期処理内ですべて呼び出すとそれだけ処理時間がかかるだけでなく、外部APIの影響を受けることが懸念されました。そこで、外部APIの障害や遅延に備えることができ、ニアリアルタイムで処理することができる、Kafkaを利用したメッセージングシステムを利用することにしました。
本記事では、Spring Boot と Apache Kafka を組み合わせたイベント駆動アーキテクチャの導入事例の中から、1つのトピックで複数のイベントを複数のConsumerで処理する設計パターンについてご紹介いたします。
1. なぜ1つのトピックで複数のConsumerを使うのか
課題:複雑なビジネスプロセスの処理
注文処理のような複雑なビジネスプロセスでは、以下のような複数のステップが必要になります:
注文受付 → 決済処理 → 領収書作成 → メール通知
これを1つのConsumerで処理すると以下のような問題が考えられます。
- 単一障害点になる — 1つの処理が失敗すると全体が止まります
- スケーラビリティが低い — 負荷に応じて個別にスケールできません
- 責任が集中する — コードが肥大化して保守性が下がります
設計の選択肢:トピック分離 vs 1トピック複数イベント
複数の処理を扱う場合、2つの設計アプローチがあります:
アプローチA:イベントごとにトピックを分ける
[order-placed-topic] → [OrderPlacedListener]
[payment-completed-topic] → [PaymentListener]
[receipt-created-topic] → [ReceiptListener]
課題:
- トピック数が増えて管理が煩雑になる
- 関連するイベントが分散して全体像が把握しづらい
- 新しいイベント追加時にトピック作成が必要
アプローチB:1つのトピックで複数のイベントを扱う
[order-event トピック]
├─ ORDER_PLACED イベント
├─ PAYMENT_COMPLETED イベント
└─ RECEIPT_CREATED イベント
メリット:
- 関連するイベントを1箇所に集約
- トピック管理がシンプル
- イベント追加時はヘッダーで識別(トピック作成不要)
- ビジネスプロセス全体の可視性が向上
解決策:1つのトピック + 複数のConsumerで役割分担
このシステムではイベント追加やトピック管理のシンプルさを重視し、アプローチBを採用しました。
[order-event トピック]
├─→ [PaymentListener] 決済処理専用
├─→ [ReceiptListener] 領収書作成専用
├─→ [NotificationListener] 通知専用
└─→ [OrderStatusListener] ステータス管理(DB更新)専用
メリット:
- マイクロサービス的な責任分離
- 各Consumerを独立してスケール可能
- 障害の影響範囲を限定できる
2. Consumer Groupの役割と処理の分散
KafkaのConsumer Groupは、Consumerの振る舞いを決定する重要な概念です。
パターン1:異なるConsumer Group(ブロードキャスト型)
// 決済処理用Consumer
@KafkaListener(
topics = "order-event",
groupId = "payment-group" // グループA
)
public void handlePayment(Order order) {
paymentService.process(order);
}
// メール送信用Consumer
@KafkaListener(
topics = "order-event",
groupId = "notification-group" // グループB(別グループ)
)
public void handleNotification(Order order) {
mailService.send(order);
}
動作:
[order-event: メッセージ1]
├─→ payment-group のConsumer が処理
└─→ order-status-update-group のConsumer も処理(同じメッセージ)
- 各グループが同じメッセージを受信する
- 並行して独立に処理が進む
- 処理の順序保証はない(決済とステータス管理が同時に走る)
パターン2:同じConsumer Group(負荷分散型)
// Consumer 1
@KafkaListener(
topics = "order-event",
groupId = "payment-group"
)
public void handlePayment1(Order order) { ... }
// Consumer 2(同じgroupId)
@KafkaListener(
topics = "order-event",
groupId = "payment-group"
)
public void handlePayment2(Order order) { ... }
動作:
[order-event: メッセージ1] → Consumer 1 が処理
[order-event: メッセージ2] → Consumer 2 が処理
[order-event: メッセージ3] → Consumer 1 が処理
- 同じグループ内ではメッセージを分担する(パーティションごとに割り当て)
- 処理能力を水平スケールできる
- 同じメッセージを複数Consumerが受け取ることはない
このシステムでは複数のConsumerで役割分担するため、主にパターン1を重視しています。
3. 1つのトピックで複数のイベントを扱う:event_typeヘッダーによる識別
1つのトピックで複数のイベントを扱うためには、各メッセージがどのイベントなのかを識別する仕組みが必要です。
そこで、event_typeヘッダーを使ってイベント種別を付与し、Consumer側で処理を振り分けます。
この設計により、関連するイベント(ORDER_PLACED, PAYMENT_COMPLETED, RECEIPT_CREATEDなど)を
同じorder-eventトピックに集約しながら、各Consumerが担当するイベントだけを処理できます。
Producer側:イベントタイプを設定
@Service
public class OrderEventPublisher {
public void publishOrderPlaced(Order order) {
var message = MessageBuilder
.withPayload(order)
.setHeader(KafkaHeaders.TOPIC, "order-event")
.setHeader("event_type", "ORDER_PLACED") // イベント種別
.build();
kafkaTemplate.send(message);
}
public void publishPaymentCompleted(Order order) {
var message = MessageBuilder
.withPayload(order)
.setHeader(KafkaHeaders.TOPIC, "order-event")
.setHeader("event_type", "PAYMENT_COMPLETED") // 別のイベント種別
.build();
kafkaTemplate.send(message);
}
}
Consumer側:ヘッダーで処理を振り分け
@Component
@RequiredArgsConstructor
public class OrderStatusListener {
private final OrderStatusUpdateService orderStatusUpdateService;
@KafkaListener(
topics = "order-event",
groupId = "order-status-update-group"
)
public void handleOrderEvent(
@Payload Order order,
@Header(name = "event_type", required = false) String eventType,
Acknowledgment acknowledgment
) {
log.info("受信: orderId={}, eventType={}", order.getId(), eventType);
try {
// イベントタイプに応じて処理を振り分け
switch (eventType) {
case "ORDER_PLACED":
// 受注完了時のステータス更新
orderStatusUpdateService.updateStatusToOrderPlaced(order);
break;
case "PAYMENT_COMPLETED":
// 決済完了時のステータス更新
orderStatusUpdateService.updateStatusToPaymentCompleted(order);
break;
case "RECEIPT_CREATED":
// 領収書作成完了時のステータス更新
orderStatusUpdateService.updateStatusToReceiptCreated(order);
break;
default:
log.debug("対象外イベント: {}", eventType);
}
acknowledgment.acknowledge();
} catch (Exception e) {
log.error("ステータス更新失敗: orderId={}, eventType={}",
order.getId(), eventType, e);
throw e; // リトライ
}
}
}
ポイント:
- 1つのListenerで複数のイベントタイプを処理
- イベントドリブンなステートマシンとして機能
- 処理ロジックを集約できる
このシステムでは外部APIの呼び出しとステータス管理(DB更新)で責務を分離したため、ステータス管理Consumerではevent_typeにより処理を振り分けています。
4. イベント駆動アーキテクチャ:並行処理と順序処理の実現
1つのトピックで複数のイベントを扱うことで、並行処理と順序処理の両方を実現できます。
各Consumerが処理完了後に次のイベントを発行することで、イベントチェーンによる順序処理を行いながら、異なるConsumer Group間では並行処理が可能になります。
アーキテクチャ:複数の異なるイベントが同じトピックを流れる
┌─────────────┐
│ Controller │ POST /orders
└──────┬──────┘
↓ 【ORDER_PLACED】イベント発行
┌──────────────────────────────────────────────┐
│ Kafka: order-event トピック(単一トピック) │
│ ├─ ORDER_PLACED イベント │
│ ├─ PAYMENT_COMPLETED イベント │
│ └─ RECEIPT_CREATED イベント │
└──────┬───────────────────────────────────────┘
↓ すべてのConsumerが購読(異なるConsumer Group)
|
├──→ ┌────────────────────┐
│ │ PaymentListener │ 決済処理(ORDER_PLACEDを処理)
│ │ (groupId: A) │
│ └──────┬─────────────┘
│ ↓ 【PAYMENT_COMPLETED】発行(同じトピックへ)
│
├──→ ┌────────────────────┐
│ │ReceiptListener │ 領収書作成(PAYMENT_COMPLETEDを処理)
│ │ (groupId: B) │
│ └──────┬─────────────┘
│ ↓ 【RECEIPT_CREATED】発行(同じトピックへ)
│
├──→ ┌────────────────────┐
│ │NotificationListener│ メール送信(RECEIPT_CREATEDを処理)
│ │ (groupId: C) │
│ └────────────────────┘
│
└──→ ┌────────────────────┐
│OrderStatusListener │ ステータス更新(全イベントを処理)
│ (groupId: D) │
└────────────────────┘
【重要】すべてのイベントが order-event という1つのトピックを経由
実装例
1. PaymentListener(決済処理)
@Component
@RequiredArgsConstructor
public class PaymentListener {
private final PaymentService paymentService;
private final OrderEventPublisher eventPublisher;
@KafkaListener(
topics = "order-event",
groupId = "payment-group"
)
public void handleOrderPlaced(
@Payload Order order,
@Header("event_type") String eventType,
Acknowledgment ack
) {
if (!"ORDER_PLACED".equals(eventType)) {
ack.acknowledge(); // 対象外イベントはスキップ
return;
}
try {
// 決済処理実行
paymentService.processPayment(order);
// 次のイベントを発行
eventPublisher.publishPaymentCompleted(order);
ack.acknowledge();
log.info("決済完了: orderId={}", order.getId());
} catch (Exception e) {
log.error("決済失敗: orderId={}", order.getId(), e);
throw e; // リトライ
}
}
}
2. ReceiptListener(領収書作成)
@Component
@RequiredArgsConstructor
public class ReceiptListener {
private final ReceiptService receiptService;
private final OrderEventPublisher eventPublisher;
@KafkaListener(
topics = "order-event",
groupId = "receipt-group"
)
public void handlePaymentCompleted(
@Payload Order order,
@Header("event_type") String eventType,
Acknowledgment ack
) {
if (!"PAYMENT_COMPLETED".equals(eventType)) {
ack.acknowledge();
return;
}
try {
// 領収書作成処理
receiptService.createReceipt(order);
// 次のイベントを発行
eventPublisher.publishReceiptCreated(order);
ack.acknowledge();
log.info("領収書作成完了: orderId={}", order.getId());
} catch (Exception e) {
log.error("領収書作成失敗: orderId={}", order.getId(), e);
throw e;
}
}
}
フロー図
時刻 t0: ORDER_PLACED イベント発行
↓
├─────────────────────────────┐
↓ ↓
時刻 t1: PaymentListener が処理 OrderStatusListener が処理
決済API呼び出し DB更新(受注完了)
↓
PAYMENT_COMPLETED イベント発行
↓
├─────────────────────────────┐
↓ ↓
時刻 t2: ReceiptListener が処理 OrderStatusListener が処理
領収書作成API呼び出し DB更新(決済完了)
↓
RECEIPT_CREATED イベント発行
↓
├─────────────────────────────┐
↓ ↓
時刻 t3: NotificationListener OrderStatusListener が処理
メール送信API呼び出し DB更新(領収書作成完了)
重要な特性:
- 各Listenerは異なるConsumer Groupに所属→同じイベントを複数のConsumerが受信
- event_typeで処理を振り分け→必要なイベントだけを処理
- 外部API呼び出しは順序処理(イベントチェーン)→処理完了後に次のイベントを発行
- ステータス管理(DB更新)は並行処理→全てのイベントを並行して処理
- 各ステージが独立してリトライ・スケール可能
5. 実装パターン
実装パターンをまとめると以下の2パターンになります。
パターンA:1 Listener 1 Event Type(推奨)
// 決済専用Listener
@KafkaListener(topics = "order-event", groupId = "payment-group")
public void handleOrderPlaced(...) {
if (!"ORDER_PLACED".equals(eventType)) return;
// 決済処理のみ
}
// 領収書専用Listener
@KafkaListener(topics = "order-event", groupId = "receipt-group")
public void handlePaymentCompleted(...) {
if (!"PAYMENT_COMPLETED".equals(eventType)) return;
// 領収書作成処理のみ
}
メリット:
- 責任が明確
- 独立したスケーリングが容易
- コードの可読性が高い
デメリット:
- Listenerクラスが増える
- 対象外イベントもConsumerに届く(フィルタ処理が必要)
パターンB:1 Listener Multiple Event Types
@KafkaListener(topics = "order-event", groupId = "order-status-update-group")
public void handleAllOrderEvents(
@Payload Order order,
@Header("event_type") String eventType,
Acknowledgment ack
) {
switch (eventType) {
case "ORDER_PLACED": updateStatusToOrderPlaced(order); break;
case "PAYMENT_COMPLETED": updateStatusToPaymentCompleted(order); break;
case "RECEIPT_CREATED": updateStatusToReceiptCreated(order); break;
}
ack.acknowledge();
}
メリット:
- コードが集約される
- イベントの流れが1箇所で把握できる
デメリット:
- クラスが肥大化しやすい
- 個別のスケーリングができない(Consumer Group が1つ)
- テストが複雑になる
このシステムでは基本的にパターンAで実装していますが、ステータス管理(DB更新)はコード集約するためにパターンBで実装しています。
6. まとめ:1つのトピック + 複数イベント + 複数Consumer = 柔軟なイベント駆動設計
この設計パターンの核心
本記事でご紹介した設計の核心は、**「関連する複数のイベント(ORDER_PLACED, PAYMENT_COMPLETED, RECEIPT_CREATEDなど)を1つのトピック(order-event)に集約し、複数のConsumerがそれぞれ担当イベントを処理する」**という点です。
これにより、トピックの乱立を避けながら、マイクロサービス的な責任分離とスケーラビリティを実現できます。
この設計パターンのメリット
| 要素 | 説明 |
|---|---|
| 責任分離 | 各Consumerが特定の処理だけを担当(Single Responsibility) |
| 独立スケーリング | 負荷に応じて個別にConsumer数を増やせる |
| 障害分離 | 1つのConsumerが落ちても他は処理を続行 |
| 柔軟性 | 新しい処理を追加する際、既存Consumerに影響なし |
| 並行・順序の両立 | Consumer Groupの分離により、並行処理と順序処理を同時に実現 |
この設計パターンを実現する3つの要素
- Consumer Groupの分離 — 異なるGroupで同じイベントを受信、同じGroupでメッセージを分担
- event_typeヘッダー — Consumerが自分の担当イベントを判定して処理を振り分け
- イベントチェーン(オプション) — 順序処理が必要な場合、処理完了後に次のイベントを発行
設計のポイント
【良い設計】
- 各Consumerが1つの責任を持つ(Single Responsibility Principle)
- イベント駆動で疎結合なアーキテクチャ
- Consumer Groupを適切に分離して並行処理と順序処理を使い分ける
- event_typeヘッダーで柔軟な処理の振り分けを実現
【避けるべき設計】
- 1つのConsumerで全部やる → スケールしない、障害の影響範囲が広い
- 本来並行処理すべき処理を順序処理にする → スループットが低下
- Consumer Group を不適切に共有する → 必要なイベントが届かない
アーキテクチャの特徴
このシステムでは、以下のように設計しています:
順序処理が必要な処理(外部API呼び出し)
- 各Consumerが異なるConsumer Groupに所属
- イベントチェーンで順序を保証
- 例:決済処理 → 領収書作成 → メール送信
並行処理で良い処理(ステータス管理・DB更新)
- 全てのイベントを受信して並行処理
- event_typeで処理を振り分け
- 外部API呼び出しと同時に実行可能
このパターンを使うことで、複雑なビジネスプロセスを、拡張性と保守性を保ちながら、効率的に実装できます。