Debezium + Amazon MSK でアウトボックスパターンによるサービス間連携を構築してみた

Debezium + Amazon MSK でアウトボックスパターンによるサービス間連携を構築してみた

Spring Boot + JPA でアウトボックスパターンを実装し、Debezium を使ってイベントを Amazon MSK に CloudEvents 1.0 形式で送信する方法を試してみました。
2026.07.30

AWS 上でマイクロサービス間のイベント連携を行う場合、Kinesis や EventBridge、Kafka を利用するなど様々なパターンがあるかと思います。
今回は調査も兼ねて、Spring Boot + JPA を使ったサービス間連携の実装方法としてアウトボックスパターンを採用し、Debezium を使って Amazon MSK (Managed Streaming for Apache Kafka) にイベントを CloudEvents 1.0 形式で送信する方法を試してみます。

  • 技術スタックの構成
    • App: Spring Boot(Kotlin) + JPA
    • DB: Aurora MySQL
    • Connector: Debezium
    • Queue: Amazon MSK (Managed Streaming for Apache Kafka)
    • Event Schema: CloudEvents 1.0

outbox-pattern-aws

サンプルコード

GitHub にサンプルコードを置いていますので、興味がある方はそちらを参照ください。あくまでデモ用のため、パスワードなどの機密情報は平文で書いている点はご注意ください。

https://github.com/seiichi1101/outbox-pattern-in-aws

アプリケーション(アウトボックスパターン)

マイクロサービス間の連携では、「DB を更新し、あわせてイベントも発行する」という処理が必要になります。このとき DB とメッセージブローカーという 2 つの異なるシステムへ書き込むことになるため、片方だけが成功して片方が失敗する、いわゆる dual write 問題 が発生します。DB 更新後・イベント送信前にアプリケーションがクラッシュするとイベントが失われますし、逆順にすれば今度は DB に存在しないデータのイベントが流れてしまいます。

アウトボックスパターンでは、outbox_events テーブルを用意しておき、業務データの更新と同じ DB トランザクション内でイベントを outbox_events テーブルにも書き込みます。DB の更新とイベントの記録が同一トランザクションで行われるため、DB の更新が成功した場合のみイベントが記録され、ロールバックが発生した場合はイベントの書き込みも取り消されます。これにより、DB の状態と発行されるイベントの整合性が保たれます。

outbox_events テーブルに書かれたイベントを実際にブローカーへ送信するのは、アプリケーションとは別のプロセスの仕事です。一般的には、スケジュール実行するバッチ処理(ポーリング)や、DB の変更ログを監視する CDC (Change Data Capture) ツールを使います。今回は後述する Debezium に任せます。

  • Spring Boot(Kotlin) + JPA での実装例
@Service
class UserService(
    private val userRepository: UserRepository,
    private val outboxEventRepository: OutboxEventRepository,
    private val objectMapper: ObjectMapper,
) {
    fun listUsers(): List<User> {
        return userRepository.findAll()
    }

    @Transactional
    fun register(name: String): User {
        val user = userRepository.save(User(name = name))
        saveOutboxEvent("USER_REGISTERED", user)
        return user
    }

    @Transactional
    fun update(id: Long, name: String): User {
        val user = userRepository.findById(id)
            .orElseThrow { NoSuchElementException("User not found: id=$id") }
        user.name = name
        user.updatedAt = Instant.now()
        saveOutboxEvent("USER_UPDATED", user)
        return user
    }

    @Transactional
    fun delete(id: Long) {
        val user = userRepository.findById(id)
            .orElseThrow { NoSuchElementException("User not found: id=$id") }
        userRepository.delete(user)
        saveOutboxEvent("USER_DELETED", user)
    }

    // User の登録・更新・削除のたびに、イベントを outbox_events テーブルに書き込む
    private fun saveOutboxEvent(type: String, user: User) {
        val data = objectMapper.writeValueAsBytes(
            mapOf(
                "id" to user.id,
                "name" to user.name,
            )
        )
        // outbox_events の主キーと CloudEvents の id に同じ UUID を使う
        val eventId = UUID.randomUUID().toString()
        // eventのpayloadには CloudEvents SDK for Java を利用
        // https://cloudevents.github.io/sdk-java/spring.html
        val cloudEvent = CloudEventBuilder.v1()
            .withId(eventId)
            .withSource(URI.create("/users"))
            .withType(type)
            .withSubject(user.id.toString())
            .withTime(OffsetDateTime.now(ZoneOffset.UTC))
            .withDataContentType("application/json")
            .withData(data)
            .build()
        outboxEventRepository.save(
            OutboxEvent(
                id = eventId,
                aggregateType = "USER",
                aggregateId = user.id.toString(),
                type = type,
                payload = jsonFormat.serialize(cloudEvent).toString(Charsets.UTF_8),
            )
        )
    }

    companion object {
        private val jsonFormat = JsonFormat()
    }
}

ポイントは、アプリケーションが Kafka に直接メッセージを送信していないことです。アプリケーションは自身の DB に書き込むだけで、Kafka への送信は非同期で動作する別のプロセス(今回は Debezium)に移譲します。そのため、アプリケーション側には Kafka クライアントの設定も送信リトライの実装も不要で、DB の更新処理に集中できます。

outbox_events テーブルに対応するエンティティは以下のとおりです。

@Entity
@Table(name = "outbox_events")
class OutboxEvent(
    @Id
    @Column(columnDefinition = "char(36)")
    var id: String = UUID.randomUUID().toString(),

    @Column(name = "aggregatetype", nullable = false)
    var aggregateType: String,

    @Column(name = "aggregateid", nullable = false)
    var aggregateId: String,

    @Column(name = "type", nullable = false)
    var type: String,

    @Column(nullable = false, columnDefinition = "json")
    var payload: String,

    @Column(nullable = false)
    var createdAt: Instant = Instant.now(),
)

AWS 環境では、App は ECS のサービスとして動かしています。App の ECS タスクは Aurora MySQL に接続して、ユーザーの登録・更新・削除を行います。

resource "aws_ecs_task_definition" "app" {
  family                   = "${var.name}-app"
  requires_compatibilities = ["FARGATE"]
  network_mode             = "awsvpc"
  cpu                      = 512
  memory                   = 1024
  execution_role_arn       = aws_iam_role.task_execution.arn

  container_definitions = jsonencode([
    {
      name         = "app"
      image        = "${aws_ecr_repository.app.repository_url}:${var.app_image_tag}"
      essential    = true
      portMappings = [{ containerPort = 8080 }]
      environment = [
        { name = "SPRING_DATASOURCE_URL", value = "jdbc:mysql://${aws_rds_cluster.aurora.endpoint}:3306/${var.db_name}" },
        { name = "SPRING_DATASOURCE_USERNAME", value = var.db_username },
        { name = "SPRING_DATASOURCE_PASSWORD", value = var.db_password },
        { name = "SPRING_KAFKA_BOOTSTRAP_SERVERS", value = aws_msk_cluster.main.bootstrap_brokers },
      ]
      logConfiguration = {
        logDriver = "awslogs"
        options = {
          awslogs-group         = aws_cloudwatch_log_group.app.name
          awslogs-region        = var.region
          awslogs-stream-prefix = "app"
        }
      }
    }
  ])
}

Amazon Aurora MySQL

RDBMS は特にこだわりはないですが、今回は使い慣れている Amazon Aurora MySQL を使いました。Debezium は MySQL の binlog(バイナリログ)を監視することで、DB の変更をキャプチャできます。

  • Amazon Aurora MySQL の setup
resource "aws_db_subnet_group" "aurora" {
  name       = "${var.name}-aurora"
  subnet_ids = aws_subnet.private[*].id
}

resource "aws_rds_cluster_parameter_group" "aurora" {
  name   = "${var.name}-aurora"
  family = "aurora-mysql8.0"

  # Row-based binlog is required by Debezium
  parameter {
    name         = "binlog_format"
    value        = "ROW"
    apply_method = "pending-reboot"
  }

  parameter {
    name         = "binlog_row_image"
    value        = "FULL"
    apply_method = "immediate"
  }
}

resource "aws_rds_cluster" "aurora" {
  cluster_identifier              = "${var.name}-aurora"
  engine                          = "aurora-mysql"
  database_name                   = var.db_name
  master_username                 = var.db_username
  master_password                 = var.db_password
  db_subnet_group_name            = aws_db_subnet_group.aurora.name
  vpc_security_group_ids          = [aws_security_group.aurora.id]
  db_cluster_parameter_group_name = aws_rds_cluster_parameter_group.aurora.name
  skip_final_snapshot             = true

  serverlessv2_scaling_configuration {
    min_capacity = 0.5
    max_capacity = 2
  }
}

resource "aws_rds_cluster_instance" "aurora" {
  identifier         = "${var.name}-aurora-1"
  cluster_identifier = aws_rds_cluster.aurora.id
  engine             = aws_rds_cluster.aurora.engine
  instance_class     = "db.serverless"
}

ポイントは、クラスターパラメータグループでバイナリログの形式(binlog_format)を行ベース(ROW)にしているところです。Aurora はデフォルトで binlog が実質無効(binlog_format = OFF)になっているため、これを ROW に変更することで初めて binlog が出力され、Debezium が CDC できるようになります。binlog_row_imageFULL に設定します。これにより、行ベースの binlog に各レコードの全カラム値が記録されます。UPDATE では変更前後、INSERT では変更後、DELETE では変更前の行データが記録されるため、Debezium が完全な変更イベントを生成できます。Aurora MySQL では FULL がデフォルトですが、念のため設定値を確認します。

Aurora 特有の注意点として binlog の保持期間 があります。Aurora MySQL はデフォルトでは binlog を可能な限り早く削除してしまうため、Debezium(コネクタ)が停止している間に binlog が消えると、復旧時に続きから読めなくなります。以下のストアドプロシージャで保持期間を明示的に設定しておきます。

CALL mysql.rds_set_configuration('binlog retention hours', 24);

もう 1 つの注意点が フェイルオーバー時の読み取り位置 です。Debezium の MySQL コネクタはデフォルトで binlog のファイル名とポジションで読み取り位置を管理しますが、これらはインスタンス固有の情報のため、フェイルオーバーで Writer が切り替わるとコネクタが続きを見失う可能性があります。本番運用では、クラスターパラメータグループで gtid_modeenforce_gtid_consistencyON にして GTID(グローバルトランザクション識別子)を有効化し、インスタンスを跨いで一意なトランザクション ID で読み取り位置を管理することができます。なお、今回は単一インスタンスのデモ環境のため、GTID の有効化とフェイルオーバー検証は省略しています。

https://docs.aws.amazon.com/AmazonRDS/latest/AuroraUserGuide/mysql-replication-gtid.html

outbox_events テーブルの定義は、Debezium の Outbox Event Router がデフォルトで想定している Basic Outbox Table の仕様を参考にしています。

https://debezium.io/documentation/reference/stable/transformations/outbox-event-router.html#basic-outbox-table

Column        |          Type          | Modifiers
--------------+------------------------+-----------
id            | uuid                   | not null
aggregatetype | character varying(255) | not null
aggregateid   | character varying(255) | not null
type          | character varying(255) | not null
payload       | jsonb                  |

なお、上記はドキュメント上の PostgreSQL の例なので、uuidjsonb といった型は MySQL ではそのまま使えません。今回は前述のエンティティから以下のようなテーブルが生成されます。

create table outbox_events
(
    id            char(36)     not null primary key,
    aggregateid   varchar(255) not null,
    aggregatetype varchar(255) not null,
    type          varchar(255) not null,
    payload       json         not null,
    created_at    datetime(6)  not null
);

Debezium

Debezium は、DB の変更をキャプチャして Kafka に送信する CDC (Change Data Capture) ツールです。Debezium を使うことで、DB の更新をトリガーとしてイベントを Kafka に送信することができます。Debezium は Kafka Connect のプラグイン(Source Connector)として動作するため、利用するには Kafka Connect のワーカーを立てた上でコネクタを登録する必要があります。

AWS 環境では、Debezium も ECS のサービスとして動かしています。Debezium の公式 Docker イメージ(quay.io/debezium/connect)を使って、Fargate タスクとして起動しています。なお、Kafka Connect は自身の設定・オフセット・ステータスを Kafka のトピック(connect_configs / connect_offsets / connect_statuses)に保存するため、これらのトピック名を環境変数で指定します。これらのトピックは Kafka Connect という実行基盤のメタ情報を永続化するデータストアであり、仮に Debezium のタスクがクラッシュして別のタスクが起動した場合でも、コネクタの設定やオフセットを引き継いで処理を再開できます。

resource "aws_ecs_task_definition" "connect" {
  family                   = "${var.name}-connect"
  requires_compatibilities = ["FARGATE"]
  network_mode             = "awsvpc"
  cpu                      = 1024
  memory                   = 2048
  execution_role_arn       = aws_iam_role.task_execution.arn

  container_definitions = jsonencode([
    {
      name         = "connect"
      image        = "quay.io/debezium/connect:3.6.0.Final"
      essential    = true
      portMappings = [{ containerPort = 8083 }]
      environment = [
        { name = "BOOTSTRAP_SERVERS", value = aws_msk_cluster.main.bootstrap_brokers },
        { name = "GROUP_ID", value = "outbox-connect" },
        { name = "CONFIG_STORAGE_TOPIC", value = "connect_configs" },
        { name = "OFFSET_STORAGE_TOPIC", value = "connect_offsets" },
        { name = "STATUS_STORAGE_TOPIC", value = "connect_statuses" },
      ]
      logConfiguration = {
        logDriver = "awslogs"
        options = {
          awslogs-group         = aws_cloudwatch_log_group.connect.name
          awslogs-region        = var.region
          awslogs-stream-prefix = "connect"
        }
      }
    }
  ])
}

Debezium では、単に DB の変更をキャプチャするいわゆる CDC としての使い方だけでなく、Outbox Event Router という機能を使うことで、アウトボックスパターンを簡単に実装することができます。

https://debezium.io/documentation/reference/stable/transformations/outbox-event-router.html

Outbox Event Router は、Kafka Connect の SMT (Single Message Transformation) として実装されています。SMT は、Connector がキャプチャしたレコードを Kafka へ書き込む直前に 1 メッセージずつ変換する仕組みです。通常、Debezium がキャプチャした変更イベントは変更前後の行データ(before / after)を含む独自のエンベロープ形式ですが、Outbox Event Router SMT はここから outbox_events テーブルへの INSERT を取り出し、

  • payload カラムの中身をそのままメッセージの値にする
  • aggregatetype のカラム値からルーティング先のトピック名を組み立てる
  • aggregateid のカラム値をメッセージのキーにする
  • 任意のカラムをヘッダーやエンベロープに追加する

という変換を行った上で Kafka に送信してくれます。つまり、Outbox テーブルのレコードからイベントメッセージへの変換ロジックを一切コーディングせず、設定だけで実現できます。ちなみに outbox_events テーブルへの DELETE は SMT がデフォルトで無視するため、テーブルの肥大化を防ぐために定期的にレコードを削除しても、それが余計なイベントとして流れることはありません(削除自体は binlog に載るので、Connect 側の処理量としては乗ってきます)。

タスク起動後の初期設定として、Debezium のコネクタを登録します。コネクタは Kafka Connect の REST API (POST /connectors) で作成できます。なお、今回はデモ用のため DB のユーザー名とパスワードを平文で書いていますが、Kafka Connect には設定値を外部から解決する Config Provider の仕組みがあるので、実運用ではそちら経由で Secrets Manager などから取得するとよいでしょう。

{
  "name": "outbox-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.hostname": "REPLACED_BY_AURORA_ENDPOINT",
    "database.port": "3306",
    "database.user": "root",
    "database.password": "password",
    "database.server.id": "184054",
    "topic.prefix": "outbox-app",
    "database.include.list": "test",
    "table.include.list": "test.outbox_events",
    "schema.history.internal.kafka.bootstrap.servers": "REPLACED_BY_MSK_BROKERS",
    "schema.history.internal.kafka.topic": "schema-changes.outbox",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    "transforms": "outbox,cloudEventsContentType",
    "transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
    "transforms.outbox.table.fields.additional.placement": "type:header:eventType",
    "transforms.outbox.table.expand.json.payload": "true",
    "transforms.cloudEventsContentType.type": "org.apache.kafka.connect.transforms.InsertHeader",
    "transforms.cloudEventsContentType.header": "content-type",
    "transforms.cloudEventsContentType.value.literal": "application/cloudevents+json"
  }
}

https://debezium.io/documentation/reference/stable/transformations/outbox-event-router.html#outbox-event-router-configuration-options

  • SMT の設定

今回は outbox(Outbox Event Router)と cloudEventsContentType(InsertHeader)の 2 つの SMT を適用しています。transforms.outbox.* が Outbox Event Router SMT の設定で、それぞれの役割は以下のとおりです。

設定 役割
transforms 適用する SMT のエイリアス。複数指定する場合はカンマ区切りで書き、記述した順に変換が適用される
transforms.outbox.type SMT のクラス名。Outbox Event Router の場合は io.debezium.transforms.outbox.EventRouter を指定
transforms.outbox.table.fields.additional.placement 追加で載せたいカラムの指定。今回は type カラムの値を eventType ヘッダーとして付与しており、コンシューマーはペイロードをパースせずにイベント種別でフィルタリングできる
transforms.outbox.table.expand.json.payload payload カラムの JSON 文字列を、エスケープされた文字列ではなく JSON オブジェクトとして展開して送信する

今回はテーブルのカラム名を Debezium のデフォルト値どおり(aggregatetype / aggregateid / payload)にしているため、ルーティングやキーに関する設定は省略しています。デフォルトの挙動は以下のとおりです。

  • aggregatetype カラム(route.by.field)の値がトピック名テンプレート outbox.event.${routedByValue}route.topic.replacement)に展開されるため、今回は outbox.event.USER トピックへルーティングされる
  • aggregateid カラム(table.field.event.key)がメッセージのキーになるため、aggregateid が同じユーザーのイベントは同じパーティションに書き込まれる
  • payload カラム(table.field.event.payload)がメッセージの値になる

カラム名をデフォルト値以外の名前にしている場合は、これらの設定でマッピングを上書きできます。

もう 1 つの transforms.cloudEventsContentType.* は、Kafka Connect 組み込みの InsertHeader SMT の設定で、すべてのメッセージに content-type: application/cloudevents+json ヘッダーを付与しています。後述のとおり、今回はメッセージの値に CloudEvents 形式の JSON がそのまま入る(structured mode 相当の)形になるため、その旨をヘッダーで明示する目的です。

  • スキーマ履歴トピック(MySQL コネクタ特有の設定)

設定の中に、Outbox とは直接関係のない以下の項目があります。

  "schema.history.internal.kafka.bootstrap.servers": "REPLACED_BY_MSK_BROKERS",
  "schema.history.internal.kafka.topic": "schema-changes.outbox",

これは MySQL コネクタ特有のデータベーススキーマ履歴(database schema history)の設定です。MySQL の binlog に記録されている行イベントには、その時点のテーブル定義(カラム名や型)は含まれていません。そのため Debezium は、binlog 中の DDL 文をパースして「各時点でのテーブルスキーマ」を内部的に管理しており、この履歴を専用の Kafka トピックに永続化しています。コネクタが再起動した際は、このトピックからスキーマ履歴を読み戻して、binlog の続きを正しく解釈できる状態を復元します。

このトピックはあくまでコネクタの内部用であり、アプリケーションが Subscribe するためのものではありません。

https://debezium.io/documentation/reference/stable/connectors/mysql.html#mysql-schema-history-topic

Amazon MSK (Managed Streaming for Apache Kafka)

Amazon MSK は、AWS が提供するマネージドな Kafka サービスです。Kafka クラスターのセットアップやスケーリングといった運用面を AWS が代行してくれるため、Kafka の運用負荷を削減することができます。

  • Amazon MSK の setup
resource "aws_msk_configuration" "main" {
  name           = "${var.name}-msk"
  kafka_versions = ["3.9.x"]

  # Debezium and Kafka Connect create topics on demand
  server_properties = <<-PROPERTIES
    auto.create.topics.enable=true
    default.replication.factor=3
    min.insync.replicas=2
  PROPERTIES
}

resource "aws_msk_cluster" "main" {
  cluster_name           = "${var.name}-msk"
  kafka_version          = "3.9.x"
  number_of_broker_nodes = 3

  broker_node_group_info {
    instance_type   = "kafka.t3.small"
    client_subnets  = aws_subnet.private[*].id
    security_groups = [aws_security_group.msk.id]
    storage_info {
      ebs_storage_info {
        volume_size = 20
      }
    }
  }

  configuration_info {
    arn      = aws_msk_configuration.main.arn
    revision = aws_msk_configuration.main.latest_revision
  }

  # Plaintext inside the VPC to keep app / connect config simple (demo only)
  encryption_info {
    encryption_in_transit {
      client_broker = "PLAINTEXT"
      in_cluster    = true
    }
  }
}

ポイントは auto.create.topics.enable=true です。これは、存在しないトピックへ書き込みがあったときに Kafka が自動でトピックを作成する設定です。MSK のデフォルトは false のため、これを有効にしておかないと、Debezium の Outbox Event Router が最初のイベントをトピックへ書き込もうとした時点で「トピックが存在しない」というエラーで失敗します。今回の構成では事前にトピックを作成する仕組みを用意していないため、必要な設定です。

なお、自動作成されるトピックはブローカーのデフォルト設定(パーティション数 num.partitions=1 など)で作られます。本番用途でパーティション数やリテンションを制御したい場合は、自動作成に頼らず、事前にトピックを明示的に作成しておくほうがよいでしょう。

順序保証と配信保証

順序保証はパーティション単位

Kafka がメッセージの順序を保証するのは、トピックのパーティション内のみです。同一のキーを持つメッセージは(デフォルトのパーティショナであれば)同一のパーティションに書き込まれるため、キー単位では送信された順に受け取れます。Outbox Event Router は table.field.event.key に指定したカラム(今回はデフォルトの aggregateid)をキーとして使うため、同一ユーザーのイベントは必ず同じパーティションに書き込まれ、コンシューマーは「ユーザー単位では」正しい順で受け取れます。逆に言うと、異なるユーザー間のイベントは、キーが異なりパーティションをまたぎうるため順序が保証されません。このように集約 ID をキーにすることで、「集約単位の順序保証」を確保できます。

配信保証は at-least-once

この構成のエンドツーエンドの配信保証は at-least-once です。Debezium (Kafka Connect) は binlog の読み取り位置(オフセット)を定期的に Kafka へコミットしますが、クラッシュや再起動のタイミングによっては、コミット前に処理したイベントを復旧後にもう一度送信することがあります。つまり、イベントが失われることはないが、重複して届くことはある という前提でコンシューマーを設計する必要があります。

そのため、コンシューマー側の処理を冪等にするのが基本です。例えば、処理済みイベントの id を記録しておいて重複を読み飛ばす、といった重複排除の仕組みを入れることが考えられます。EOS(exactly.once.source.support)を利用して、Kafka への書き込みを exactly-once にするといった方法もありますが、コンシューマーが「処理」と「オフセットコミット」をアトミックに行えなければエンドツーエンドでは重複しえます。結局のところ、コンシューマー側を冪等に作るのが現実的な解という結論になります。

CloudEvents 1.0

CloudEvents は CNCF が策定しているイベントデータの標準仕様で、イベントの ID・発生元・種別といったメタデータの持ち方が統一されるため、異なるサービス間でのイベント連携が容易になります。

https://github.com/cloudevents/spec/blob/v1.0.2/cloudevents/spec.md

先ほどのアウトボックスパターンの実装例では、outbox_events テーブルの payload カラムに CloudEvents 1.0 形式のイベントを保存していました。受信側のサービスでは、Kafka から受け取ったメッセージを CloudEvents としてパースして処理します。以下は Spring Kafka を使った受信側の実装例です。

@Service
class KafkaEventService {
...
    @KafkaListener(topics = ["outbox.event.USER"], groupId = "consumer-group-id")
    fun onMessage(record: ConsumerRecord<String, String>, @Header("eventType", required = false) eventType: String?) {
        // eventType ヘッダーを使って、イベント種別をログに出力
        log.info("Received event: {}", eventType)
        val event = try {
            // CloudEvents SDK for Java の JsonFormat でパース
            jsonFormat.deserialize(record.value().toByteArray())
        } catch (e: Exception) {
            log.warn("Skipping non-CloudEvent message at offset {}: {}", record.offset(), e.message)
            return
        }
        broadcast(event)
    }
}

実際に Consumer 側で受信するイベントは下記のような形式になります。

{
  "specversion": "1.0",
  "id": "f1bfa6e0-173c-4a31-b3e0-66439bd90468",
  "source": "/users",
  "type": "USER_REGISTERED",
  "datacontenttype": "application/json",
  "subject": "4",
  "time": "2026-07-30T08:45:35.263905545Z",
  "data": {
    "id": 4,
    "name": "Seiichi Arai"
  }
}

今回は outbox テーブルの payload カラムにイベント全体を保存しているため Kafka Protocol Binding - structured mode 相当の形になります。これは、イベント全体(メタデータ + データ)を 1 つの JSON としてメッセージの値に入れるため、content-typeapplication/cloudevents+json になります。

各言語向けに SDK が提供されており、Java では CloudEvents SDK for Java を使うことで、CloudEvents 形式のイベントの生成やパースを簡単に行うことができます。

https://cloudevents.github.io/sdk-java/spring.html

まとめ

いかがだったでしょうか。
今回は、Spring Boot + JPA を使ったサービス間連携の実装方法としてアウトボックスパターンを採用し、Debezium を使って Amazon MSK (Managed Streaming for Apache Kafka) にイベントを CloudEvents 1.0 形式で送信する方法を試してみました。マイクロサービス間のイベント連携としては割と王道なのではないかと思います。

アプリケーション側の実装は「同一トランザクションで outbox_events テーブルにイベントを書き込む」だけで、Kafka への送信処理を一切書く必要がありません。メッセージ送信のリトライや DB 更新とメッセージ送信の不整合といった悩ましい問題を Debezium 側に寄せられるため、アプリケーションのコードは非常にシンプルに保てます。Outbox Event Router もコネクタの設定 (SMT) を書くだけでトピックへのルーティングやヘッダーへのフィールド展開まで面倒を見てくれるため、想像していたより簡単に実装できました。

一方で、インフラ構成としては MSK のブローカーに加えて、Kafka Connect (Debezium) 用の ECS タスクを常時起動しておく必要があるため、システムに導入するにはややヘビーです。今回のような検証用の小規模な構成でも、それなりのランニングコストがかかります。イベントの流量が少ないうちは、outbox テーブルをポーリングして EventBridge や SNS/SQS に転送するようなシンプルな構成でも十分かもしれません。一方で、すでに Kafka を運用している場合や、イベント量が多くスループットが求められる場合に、費用対効果が見合ってくる構成だと感じました。

また、アウトボックスパターン + Debezium で「イベントを確実に届ける」ところまでは面倒を見てもらえますが、前述のとおり配信保証は at-least-once なので、重複への対処(冪等性)はコンシューマー側の責務として残ります。「アウトボックスパターンを入れれば整合性の問題がすべて解決する」わけではなく、受信側の設計とセットで考える必要がある点は意識しておきたいところです。

以上、どなたかの役に立てば幸いです。お疲れ様でした!

この記事をシェアする

関連記事