Pub/Sub は、Publisher が Topic という名前付きの配信経路にメッセージを 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 API | APIリクエストJSON | クライアントが注文APIに送る注文内容 |
Order API -> Order DB | 注文データ | 注文APIが生成してDBに保存する正本 |
Order API -> Topic | Pub/Subメッセージ | 注文APIが後続処理へ伝えるイベント |
Order API -> Client | APIレスポンス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,
});
この指定によって、message は order-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_id や user_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-sub や update-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-sub や email-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-sub と email-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 を使った処理設計を読みやすくなります。