Apache Kafkaの基礎からSpring Boot連携実装まで

メッセージキュー導入の背景と利点

現代の分散システムにおいて、メッセージキュー(MQ)は欠かせないコンポーネントです。MQを導入することで、システム間の結合度を低減(デカップリング)し、非同期処理によるレスポンス時間の短縮や、急激なアクセス増加(スパイク)に対する負荷平準化(シェーピング)が可能になります。

一方で、システム構成が複雑化したり、データの一貫性維持が難しくなるというデメリットも存在します。これらの課題を克服しつつ、大規模なデータストリーミングを実現する代表的なプラットフォームとして「Apache Kafka」が注目されています。

Kafkaのコア概念とアーキテクチャ

Kafkaは単なるメッセージキューではなく、高スループットな分散ストリーミングプラットフォームです。その特徴は以下の通りです。

  • 高スループット: ディスクへの順次書き込みを最適化しており、メモリへのランダムアクセスよりも高速な処理を実現します。
  • 永続化: メッセージをディスクに保存し、障害発生時のデータ損失を防ぎます。
  • 高可用性と拡張性: 分散クラスタ構成により、単一障害点(SPOF)を排除し、サーバーの増設によるスケールアウトが容易です。

主要な用語

  • Broker: Kafkaクラスタを構成する各サーバー节点。
  • Topic: メッセージの論理的な分類名(発行者と購読者の間のチャネル)。
  • Partition: Topicを物理的に分割した単位。各パーティション内ではメッセージが順序通り保存され、Offset(インデックス)によって管理されます。
  • Replica: データの複製。Leader Replicaが読み書きを担当し、Follower Replicaはデータのバックアップを行います。Leaderに障害が発生した場合、Followerのいずれかが新たなLeaderに昇格します。

ローカル環境での構築と動作確認

ここではKafkaの基本的な動作をコマンドラインで確認します。なお、Kafkaの実行にはZooKeeperが必要です(※最近のバージョンではKRaftモードでZooKeeperなしでも動作可能ですが、ここでは従来の構成で解説します)。

1. サービスの起動

まず、設定ファイルを元にZooKeeperを起動し、続いてKafka Brokerを起動します。

# ZooKeeperの起動
bin/zookeeper-server-start.sh config/zookeeper.properties

# Kafka Brokerの起動
bin/kafka-server-start.sh config/server.properties

2. Topicの作成と確認

検証用のTopicを作成します。ここではパーティション数1、レプリケーション係数1でquick-start-eventsというTopicを作成しています。

bin/kafka-topics.sh --create --bootstrap-server localhost:9092 \
  --replication-factor 1 --partitions 1 --topic quick-start-events

# Topic一覧の確認
bin/kafka-topics.sh --list --bootstrap-server localhost:9092

3. メッセージの送受信

コンソールプロデューサーでメッセージを送信し、別のコンソールコンシューマーで受信します。

# プロデューサー(メッセージ送信)
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic quick-start-events

# コンシューマー(メッセージ受信、--from-beginningで最初から読み込み)
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic quick-start-events --from-beginning

Spring Bootによるアプリケーション連携

次に、JavaアプリケーションからKafkaを利用するためのSpring Boot整合 implementationを行います。

依存関係の追加

pom.xmlspring-kafkaを追加します。

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

設定ファイル

application.propertiesでサーバーアドレスとコンシューマーグループを設定します。

spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.consumer.group-id=demo-app-group
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.enable-auto-commit=true

メッセージ送信実装

DI(依存性注入)を用いてKafkaTemplateを初期化し、指定のTopicへメッセージを送信するサービスを作成します。

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;

@Service
public class EventDispatcher {

    private final KafkaTemplate<String, String> kafkaTemplate;

    @Autowired
    public EventDispatcher(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void publishEvent(String targetTopic, String messagePayload) {
        kafkaTemplate.send(targetTopic, messagePayload);
    }
}

メッセージ受信実装

@KafkaListenerアノテーションを使用して、特定のTopicを監視し、メッセージ到着時に非同期で処理を行うコンポーネントを実装します。

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class NotificationListener {

    @KafkaListener(topics = "quick-start-events")
    public void handleIncomingMessage(ConsumerRecord<String, String> record) {
        System.out.println("受信したメッセージ: " + record.value());
        // ここにビジネスロジックを実装(例: DB更新など)
    }
}

動作確認テスト

作成したクラスをテストコードから呼び出し、メッセージングが正常に行われるか確認します。

import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import static java.lang.Thread.sleep;

@SpringBootTest
public class MessagingIntegrationTest {

    @Autowired
    private EventDispatcher eventDispatcher;

    @Test
    public void verifyMessagingFlow() throws InterruptedException {
        String targetTopic = "quick-start-events";
        
        // メッセージの送信
        eventDispatcher.publishEvent(targetTopic, "テストデータ: System Start");
        eventDispatcher.publishEvent(targetTopic, "テストデータ: Processing...");

        // 非同期処理の完了を待つために一時停止
        sleep(2000);
    }
}

タグ: ApacheKafka SpringBoot MessageQueue DistributedSystems Java

8月8日 23:31 投稿