Pub/Subのデータの流れをシーケンス図で理解する

Pub/Sub は、PublisherTopic という名前付きの配信経路にメッセージを publish し、Subscription の配信設定に従って Subscriber に届ける仕組みです。

ただ、この説明だけだと「Topic と Subscription は何を管理しているのか」「いつ処理済みになるのか」「Subscriber が落ちたらどうなるのか」が見えにくいです。この記事では、注文作成イベントを例に、Pub/Sub のデータの流れをシーケンス図で整理します。

Pub/Sub の登場人物や基本構造は『Pub/Subとは何か』で扱っています。

前提の構成

注文APIが注文を受け付けたあと、メール送信処理を非同期で実行する構成を考えます。

Client
  ↓ POST /orders
Order API
  ↓ publish
Topic: order-created

Subscription: send-email-sub

Email Worker

注文APIはメール送信処理を直接呼びません。注文作成イベントを Topic に publish し、Subscription の設定に従って Email Worker に届けられます。

シーケンス図を読む前提

正常系のシーケンスを見る前に、APIリクエスト、Pub/Subメッセージ、Topic、Message、Subscription の関係を整理します。

APIリクエストとPub/Subメッセージは別物

Pub/Sub の流れを理解するときに混乱しやすいのは、クライアントが送る APIリクエストのJSON と、Publisher が Topic に送る Pub/Subメッセージ を同じものとして見てしまうことです。

この節で出てくるデータは、後述するシーケンス図の以下の矢印に対応します。

シーケンス上の位置データの種類内容
Client -> Order APIAPIリクエストJSONクライアントが注文APIに送る注文内容
Order API -> Order DB注文データ注文APIが生成してDBに保存する正本
Order API -> TopicPub/Subメッセージ注文APIが後続処理へ伝えるイベント
Order API -> ClientAPIレスポンスJSON注文APIがクライアントへ返す受付結果

まず、Client -> Order API のAPIリクエストJSONは、たとえば以下のような形です。

POST /orders HTTP/1.1
Content-Type: application/json

{
  "user_id": "user-456",
  "items": [
    {
      "sku": "book-001",
      "quantity": 2
    }
  ],
  "billing_address": {
    "postal_code": "100-0001",
    "prefecture": "Tokyo",
    "city": "Chiyoda"
  }
}

この時点では、まだ order_id はありません。order_id は、Order API がリクエストを受け取り、注文を作成するときに生成する値です。

sequenceDiagram
  participant Client as "Client"
  participant API as "Order API"
  participant DB as "Order DB"
  participant Topic as "Topic: order-created"

  Client->>API: POST /orders
  API->>API: 入力値を検証
  API->>API: order_idを生成
  API->>DB: 注文を保存
  API->>Topic: 注文作成イベントをpublish
  API-->>Client: 201 Created

注文APIの内部では、リクエストJSONをそのまま Topic に流すのではなく、DBに保存する注文データと、Topic に publish するPub/Subメッセージを別々に作ります。

app.post("/orders", async (req, res) => {
  const orderId = generateOrderId();

  // DBに保存する注文データ。注文の正本として扱う。
  const order = {
    order_id: orderId,
    user_id: req.body.user_id,
    items: req.body.items,
    billing_address: req.body.billing_address,
    status: "created",
    created_at: new Date().toISOString(),
  };

  await orderRepository.save(order);

  // Pub/Sub に流すイベント。後続処理に必要な情報だけに絞る。
  const message = {
    event_id: generateEventId(),
    event_type: "order.created",
    order_id: order.order_id,
    user_id: order.user_id,
    item_count: order.items.length,
    created_at: order.created_at,
  };

  await pubsub.publish({
    // どの Topic に送るかは、publish 時に Publisher が指定する。
    topic: "order-created",
    data: message,
    attributes: {
      // attributes は配送・分類・スキーマ判定に使う補助情報。
      event_type: "order.created",
      schema_version: "1",
    },
  });

  res.status(201).json({
    // APIレスポンスは Pub/Sub メッセージとは別に作る。
    order_id: order.order_id,
    status: "accepted",
  });
});

この例では、Client -> Order API の注文リクエスト、Order API -> Order DB の注文データ、Order API -> Topic のPub/Subメッセージ、Order API -> Client のAPIレスポンスがそれぞれ別の形になっています。

Pub/Sub に流すメッセージは、必ずしも注文データ全体である必要はありません。後続処理が必要とする情報だけを入れる設計もできます。

publishするときにTopicが決まる

メッセージがどの Topic に属するかは、Pub/Sub がJSONの中身を見て自動で判断するわけではありません。Publisher 側の実装が、publish するときに Topic を指定します。

await pubsub.publish({
  // この指定で、message は order-created Topic に送られる。
  topic: "order-created",
  data: message,
});

この指定によって、messageorder-created という Topic に送られます。

1つのAPIが、複数の Topic に publish することもあります。

await pubsub.publish({
  // 注文作成を知らせるイベント。
  topic: "order-created",
  data: orderCreatedMessage,
});

await pubsub.publish({
  // 顧客行動の分析など、別用途のイベント。
  topic: "customer-activity-recorded",
  data: activityMessage,
});

この場合、同じ POST /orders の処理の中で、注文作成イベントと顧客行動イベントを別々の Topic に送っています。APIエンドポイントと Topic は必ず1対1になるわけではありません。

Messageの中身

Message は、Publisher が Topic に publish し、Subscriber が処理するデータの単位です。Pub/Sub の Message は、アプリケーションが扱うデータ本体と、配送や分類に使うメタデータに分けて考えると理解しやすいです。

{
  // Subscriber が処理に使うアプリケーションデータ。
  "data": {
    // 重複処理の判定に使えるイベント単位のID。
    "event_id": "evt-001",
    "event_type": "order.created",
    // 必要なら Subscriber がこのIDで注文詳細を取得する。
    "order_id": "order-123",
    "user_id": "user-456",
    "item_count": 2,
    "created_at": "2026-07-10T10:00:00+09:00"
  },
  // 配送・分類・スキーマ判定に使う補助情報。
  "attributes": {
    "event_type": "order.created",
    "schema_version": "1"
  }
}

data には Subscriber が処理に使う情報を入れます。attributes には、イベント種別やスキーマバージョンのような補助情報を入れます。

重要なのは、Message は「HTTPリクエストのJSONそのもの」ではなく、Publisher が後続処理向けに作るデータだという点です。

Subscriptionは何を表すのか

Subscription は、Topic に届いた Message を Subscriber に届けるための配信設定です。設定として見ると、たとえば以下のような情報を持ちます。

subscription: send-email-sub
topic: order-created
# Subscriber が Subscription へ取りに行く方式。
delivery_type: pull
# この秒数以内に Ack されないと再配信対象になる。
ack_deadline_seconds: 30
subscriber: Email Worker

この Subscription は、order-created Topic に届いたメッセージを、Email Worker が pull して処理するための配信設定です。

Push 型の場合は、Subscriber のHTTPエンドポイントも設定に含まれます。

subscription: send-email-sub
topic: order-created
# Pub/Sub が Subscriber のHTTPエンドポイントへ送る方式。
delivery_type: push
push_endpoint: https://example.com/pubsub/send-email
# Push では成功レスポンスが返らない場合も再配信対象になる。
ack_deadline_seconds: 30
subscriber: Email API

Subscription は単なる「subscribeするという行為」でも、ただの保存箱でもありません。Topic と Subscriber の間にある配信設定であり、どのメッセージが未処理で、どのメッセージが Ack 済みかという状態も Subscription ごとに管理されます。

正常系のシーケンス

このシーケンス図の番号は、直後の見出し番号と対応しています。たとえば、図の 1 は「1. Clientが注文APIを呼ぶ」、図の 3 は「3. Order APIがTopicにpublishする」を表します。

sequenceDiagram
  autonumber
  participant Client as "Client"
  participant API as "Order API"
  participant DB as "Order DB"
  participant Topic as "Topic: order-created"
  participant Sub as "Subscription: send-email-sub"
  participant Worker as "Email Worker"

  Client->>API: POST /orders
  API->>DB: 注文を保存
  API->>Topic: publish(order.created)
  Topic->>Sub: 未Ackメッセージとして管理
  API-->>Client: 201 Created
  Worker->>Sub: pull
  Sub-->>Worker: message
  Worker->>Worker: メール送信
  Worker->>Sub: ack
  Sub->>Sub: メッセージを処理済みにする

1. Clientが注文APIを呼ぶ

Client → Order API: POST /orders

クライアントは通常のHTTP APIを呼ぶだけです。Pub/Sub を使っているかどうかは、クライアントからは見えません。

POST /orders HTTP/1.1
Content-Type: application/json

{
  "user_id": "user-456",
  "items": [
    {
      "sku": "book-001",
      "quantity": 2
    }
  ]
}

このリクエストは、あくまで注文APIに送る入力です。このJSONがそのまま Pub/Sub に流れるわけではありません。

2. Order APIが注文を保存する

Order API → Order DB: 注文を保存

注文データそのものはDBに保存します。Pub/Sub は永続的な保存先ではないため、注文の正本を Pub/Sub に置くわけではありません。

INSERT INTO orders (
  id,
  user_id,
  status,
  created_at
) VALUES (
  'order-123',
  'user-456',
  'created',
  '2026-07-10T10:00:00+09:00'
);

ここで保存される注文データが、アプリケーション上の正本です。後続処理は必要に応じて order_id を使い、この注文データを参照します。

3. Order APIがTopicにpublishする

Order API → Topic: publish(order.created)

Order API は注文作成後に、注文作成イベントを Topic に送ります。

たとえば以下のようなメッセージです。

{
  // イベントを一意に識別するID。重複処理の判定にも使える。
  "event_id": "evt-001",
  "event_type": "order.created",
  // イベントの対象になった注文ID。
  "order_id": "order-123",
  "user_id": "user-456",
  "total_amount": 9800,
  "created_at": "2026-07-10T10:00:00+09:00"
}

ここで送っているのは「注文データの全体」ではなく、「注文が作られた」というイベントです。後続処理が詳細な注文情報を必要とする場合は、order_id を使ってDBやAPIから取得する設計もあります。

4. Subscription単位で未処理状態が管理される

Topic → Subscription: 未Ackメッセージとして管理

Topic に publish されたメッセージは、その Topic に紐づく Subscription の配信対象になります。ここで重要なのは、Subscriber が処理して Ack するまで、未処理の状態が Subscription 単位で管理されることです。Email Worker が一時的に落ちていても、メッセージはすぐに失われません。

subscription: send-email-sub
message_id: msg-001
ack_state: unacked
delivery_attempt: 0

この状態は、Subscriber が処理に成功して Ack するまで残ります。未Ackの状態があるから、Subscriber が一時的に止まっても後で再配信できます。

5. Order APIがクライアントへレスポンスを返す

Order API → Client: 201 Created

Order API は、メール送信の完了を待たずにレスポンスを返せます。これが非同期化の大きな利点です。同期処理の場合、メール送信が遅いと注文APIのレスポンスも遅くなります。

同期処理:
Order API → Email API → Clientへレスポンス

Pub/Sub を挟むと、後続処理を待たずに応答できます。

非同期処理:
Order API → Pub/Sub
Order API → Clientへレスポンス
HTTP/1.1 201 Created
Content-Type: application/json

{
  "order_id": "order-123",
  "status": "accepted"
}

このレスポンスは「注文を受け付けた」という意味です。メール送信まで完了したことを表すレスポンスではありません。

6-7. SubscriberがSubscriptionから読む

Email Worker → Subscription: pull
Subscription → Email Worker: message

Subscriber は Topic を直接読むのではなく、Subscription から読みます。Pull 型では Subscriber が Subscription に取りに行き、Push 型では Pub/Sub が Subscriber のHTTPエンドポイントにメッセージを送ります。

Pull:
Subscriber → Subscription: メッセージを取りに行く

Push:
Subscription → Subscriber: HTTPでメッセージを送る

Pull で Subscriber が受け取る情報は、アプリケーションが publish した data だけではありません。メッセージIDや属性など、Pub/Sub側の情報も一緒に受け取ります。

{
  // Pub/Sub が付与するメッセージID。アプリ側の event_id とは別物。
  "message_id": "msg-001",
  // Publisher が publish したアプリケーションデータ。
  "data": {
    "event_id": "evt-001",
    "event_type": "order.created",
    "order_id": "order-123",
    "user_id": "user-456",
    "item_count": 2,
    "created_at": "2026-07-10T10:00:00+09:00"
  },
  // Publisher が付けた補助情報。
  "attributes": {
    "event_type": "order.created",
    "schema_version": "1"
  }
}

Push の場合は、Pub/Sub が Subscriber のHTTPエンドポイントにリクエストを送ります。Google Cloud Pub/Sub の Push Subscription では、概念的には以下のようなHTTPリクエストが Subscriber に届きます。

POST /pubsub/send-email HTTP/1.1
Content-Type: application/json

{
  "message": {
    "messageId": "msg-001",
    "data": "eyJldmVudF9pZCI6ImV2dC0wMDEiLCJldmVudF90eXBlIjoib3JkZXIuY3JlYXRlZCJ9",
    "attributes": {
      "event_type": "order.created",
      "schema_version": "1"
    },
    "publishTime": "2026-07-10T01:00:00Z"
  },
  "subscription": "projects/example/subscriptions/send-email-sub"
}

この例の data はBase64エンコードされた文字列です。Subscriber 側では、この data をデコードしてアプリケーション用のJSONとして扱います。

app.post("/pubsub/send-email", async (req, res) => {
  const pubsubMessage = req.body.message;
  // Push で届く data はBase64文字列なので、まずJSON文字列へ戻す。
  const decoded = Buffer.from(pubsubMessage.data, "base64").toString("utf8");
  const event = JSON.parse(decoded);

  await sendOrderCreatedEmail({
    orderId: event.order_id,
    userId: event.user_id,
  });

  // 2xx レスポンスを返すことで、Push Subscription では Ack 相当になる。
  res.status(204).send();
});

Push 型では、Subscriber が成功レスポンスを返すことが Ack に相当します。処理に失敗したり、成功レスポンスを返せなかったりすると、Pub/Sub は後で再配信します。

8. Subscriberが処理する

Email Worker: メール送信

Subscriber は受け取ったメッセージをもとに処理します。この例では、order_iduser_id を使ってメール内容を作り、注文完了メールを送ります。

await sendOrderCreatedEmail({
  orderId: event.order_id,
  userId: event.user_id,
});

必要な情報が Pub/Sub メッセージに入っていない場合は、order_id を使ってDBやAPIから注文詳細を取得してから処理します。

9. Subscriberがackする

Email Worker → Subscription: ack

Ack は「このメッセージの処理が完了した」と Pub/Sub に伝える操作です。Pub/Sub は ack を受け取るまで、そのメッセージが正常に処理されたとは判断しません。

await message.ack();

Pull 型では、Subscriber 側の処理が成功したあとに Ack します。処理前に Ack すると、途中で失敗したときに再配信されなくなるためです。

10. Subscription上でメッセージが処理済みになる

Subscription: メッセージを処理済みにする

Ack されたメッセージは、Subscription 上で処理済みになります。以後、同じ Subscription からは通常そのメッセージは配信されません。

subscription: send-email-sub
message_id: msg-001
ack_state: acked

処理済みになるのは send-email-sub という Subscription 上での話です。別の Subscription が同じ Topic に紐づいている場合、その Subscription では別途 Ack が必要です。

Ackされなかった場合

Subscriber がメッセージを受け取ったあと、処理中に落ちた場合を考えます。

sequenceDiagram
  participant Sub as "Subscription: send-email-sub"
  participant Worker1 as "Email Worker 1"
  participant Worker2 as "Email Worker 2"

  Worker1->>Sub: pull
  Sub-->>Worker1: message
  Worker1->>Worker1: メール送信中に異常終了
  Note over Sub: ackされない
  Sub-->>Worker2: 一定時間後に再配信
  Worker2->>Worker2: メール送信
  Worker2->>Sub: ack

Ack されなかったメッセージは、一定時間後に再配信されます。

そのため、Pub/Sub では同じメッセージが複数回届く可能性があります。これを at-least-once delivery と呼びます。

「少なくとも1回は届ける」という意味で、0回にはしませんが、2回以上届く可能性はあります。

重複処理を前提にする

Pub/Sub を使う処理では、同じメッセージが2回届いても壊れないようにします。

たとえばメール送信なら、同じ event_id または order_id のメール送信履歴を保存しておき、処理済みならスキップします。

Email Worker

event_id = evt-001 は処理済みか確認
  ├─ 処理済み: スキップしてack
  └─ 未処理: メール送信して処理済みに記録してack

このように、同じ処理を複数回実行しても結果が破綻しない性質を 冪等性 と呼びます。Pub/Sub の設計では、Ack と再配信だけでなく、冪等性までセットで考えます。

ファンアウト時のデータの流れ

1つの Topic に複数 Subscription がある場合、同じメッセージが各 Subscription に届きます。

sequenceDiagram
  participant API as "Order API"
  participant Topic as "Topic: order-created"
  participant EmailSub as "Subscription: send-email-sub"
  participant AnalyticsSub as "Subscription: save-analytics-sub"
  participant InventorySub as "Subscription: update-inventory-sub"

  API->>Topic: publish(order.created)
  Topic->>EmailSub: message
  Topic->>AnalyticsSub: message
  Topic->>InventorySub: message

この場合、メール送信、分析保存、在庫更新はそれぞれ独立して進みます。send-email-sub の処理が失敗しても、save-analytics-subupdate-inventory-sub の処理には影響しません。

これが ファンアウト です。1つのイベントを複数用途に配りたいときに使います。

同じMessageでもAck状態はSubscriptionごとに分かれる

ファンアウトで混乱しやすいのは、Message 本体と Ack 状態の管理場所です。

たとえば、注文イベントを SMS 送信とメール送信の両方に使う構成を考えます。

Topic: order
  ├─ Subscription: sms-sub
  │    └─ SMS Worker
  └─ Subscription: email-sub
       └─ Email Worker

Publisher は、order Topic に1回だけ publish します。

{
  "message_id": "msg-001",
  "data": {
    "event_id": "evt-001",
    "event_type": "order.created",
    "order_id": "order-123",
    "user_id": "user-456"
  },
  "attributes": {
    "schema_version": "1"
  }
}

この Message 本体に、sms-subemail-sub の一覧が入っているわけではありません。Pub/Sub は order Topic に紐づく Subscription を見て、それぞれの Subscription でこの Message を配信対象として管理します。

以下は理解用のイメージです。実際にこの YAML がユーザーから見えるわけではありません。

topic: order
message:
  message_id: msg-001
  data:
    event_id: evt-001
    event_type: order.created
    order_id: order-123
    user_id: user-456

subscriptions:
  sms-sub:
    msg-001:
      ack_state: unacked
      delivery_attempt: 0

  email-sub:
    msg-001:
      ack_state: unacked
      delivery_attempt: 0

ここで重要なのは、Message 本体は同じでも、Ack 状態は sms-subemail-sub で別々に管理されることです。

SMS Worker が sms-sub から Message を受け取り、SMS送信に成功して Ack したとします。

await sendSms({
  orderId: "order-123",
  userId: "user-456",
});

await message.ack();

この時点で、sms-sub 側だけ処理済みになります。

subscriptions:
  sms-sub:
    msg-001:
      ack_state: acked
      delivery_attempt: 1

  email-sub:
    msg-001:
      ack_state: unacked
      delivery_attempt: 0

sms-sub では Ack 済みなので、通常は msg-001 は SMS Worker に再配信されません。一方、email-sub ではまだ Ack されていないため、メール送信は未処理のままです。

たとえば Email Worker だけが落ちている場合、SMS 側の完了には影響しません。

sms-sub:
  SMS送信成功

  ack

  msg-001 は処理済み

email-sub:
  Email Worker が down

  ack されない

  msg-001 は未処理

Push 型の email-sub であれば、Pub/Sub は Email Worker のHTTPエンドポイントへ配信を試みます。ダウン中は失敗するため Ack 扱いにならず、email-sub 側だけ再配信対象になります。

subscriptions:
  sms-sub:
    msg-001:
      ack_state: acked
      delivery_attempt: 1

  email-sub:
    msg-001:
      ack_state: unacked
      delivery_attempt: 2

Email Worker が復旧すると、email-sub に残っていた未Ackの msg-001 が再配信されます。処理に成功して 2xx レスポンスを返すと、Push Subscription では Ack 相当になります。

subscriptions:
  sms-sub:
    msg-001:
      ack_state: acked
      delivery_attempt: 1

  email-sub:
    msg-001:
      ack_state: acked
      delivery_attempt: 3

つまり、ファンアウトでは以下のように考えます。

  • Topic に publish される Message 本体は1つ
  • Topic に紐づく各 Subscription で、その Message が配信対象になる
  • Ack 状態は Message 本体ではなく、Subscription ごとに管理される
  • 片方の Subscription が Ack 済みでも、もう片方の Subscription は未Ackのまま残りうる
  • 再配信も Subscription ごとに発生する

並列ワーカー時のデータの流れ

一方、1つの Subscription を複数 Subscriber で読む場合は、同じメッセージを全員に配るわけではありません。

flowchart TB
  sub["Subscription: send-email-sub"]
  w1["Email Worker 1"]
  w2["Email Worker 2"]
  w3["Email Worker 3"]

  sub -->|"msg-1"| w1
  sub -->|"msg-2"| w2
  sub -->|"msg-3"| w3

これはファンアウトではなく、処理の分担です。メール送信処理の量が多い場合、同じ Subscription を読む Worker を増やすことで処理能力を上げられます。

まとめ

Pub/Sub のデータの流れで重要なのは以下です。

  • Publisher は Topic に publish する
  • Topic から Subscription にメッセージが渡る
  • Subscriber は Subscription から読む
  • 処理が終わったら ack する
  • ack されなければ再配信される
  • 再配信があるため、重複処理を前提にする
  • 複数 Subscription はファンアウトを表す
  • ファンアウト時の Ack 状態は Subscription ごとに管理される
  • 複数 Subscriber は処理の分担を表す

Pub/Sub は「送ったら終わり」の仕組みではありません。Subscription 単位で未Ackの状態が管理され、Subscriber が処理し、Ack によって処理済みになります。この一連の流れを理解すると、Pub/Sub を使った処理設計を読みやすくなります。

タグ: PubSub, メッセージング