Spanner 佇列情境和範例

本文提供架構模式和程式碼範例,說明如何使用 Spanner 佇列處理常見的訊息傳送情境。您可以使用這些模式,在交易提交後觸發非同步工作、安排延遲或週期性工作、透過頻外儲存空間管理大型訊息酬載、協調多事件工作流程,以及檢查或延長長期執行的背景工作的租約。

僅處理一次,且最多確認一次

「僅需處理一次」和「最多確認一次」的各種考量事項和解決方案,詳情請參閱「僅需處理一次和最多確認一次」頁面。

在交易提交後執行工作

如要在交易提交後執行工作,請在同一筆交易中將訊息傳送至佇列。

舉例來說,新使用者註冊時會觸發歡迎電子郵件:

GoogleSQL

-- Inside your application transaction:
-- 1. Insert into Users table
INSERT INTO Users (UserId, UserName) VALUES (124, 'New User');

-- 2. Send message to queue to trigger email
INSERT INTO UserTasks (UserId, MessageId, Payload)
VALUES (
  124,
  'welcome-email-id',
  b'{"type": "welcome", "email": "user@example.com"}'
);

PostgreSQL

-- Inside your application transaction:
-- 1. Insert into users table
INSERT INTO users (userid, username) VALUES (124, 'New User');

-- 2. Send message to queue to trigger email
INSERT INTO usertasks (userid, messageid, payload)
VALUES (
  124,
  'welcome-email-id',
  CAST('{"type": "welcome", "email": "user@example.com"}' AS bytea)
);

交易完成後,UserTasks 的接收者會串流傳送訊息、傳送電子郵件,並確認訊息:

GoogleSQL

-- 1. In the receiver process, stream messages from the queue
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');

-- 2. After sending the welcome email, acknowledge the message
DELETE FROM UserTasks
WHERE UserId = 124 AND MessageId = 'welcome-email-id';

PostgreSQL

-- 1. In the receiver process, stream messages from the queue
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token
FROM spanner.receive_usertasks(
    max_batch_size=>NULL, priority=>NULL, max_duration=>'20m');

-- 2. After sending the welcome email, acknowledge the message
DELETE FROM usertasks
WHERE userid = 124 AND messageid = 'welcome-email-id';

處理長時間執行的工作

如果工作可能需要比預設租約更長的時間 (超過 10 秒),請定期呼叫 SELECT * FROM RENEWLEASE_QUEUE_NAME()。

舉例來說,產生報表:

  1. 收件者會收到 RECEIVE_ReportQueue() 傳送的訊息。
  2. 開始生成報表。
  3. 每 5 秒呼叫 SELECT * FROM RENEWLEASE_ReportQueue([leaseToken]),並使用個別執行緒或常式。
  4. 完成後,請確認訊息並儲存報告。

或者,如果您有需要「最多處理一次」的長時間執行工作,或租用時間較長,請按照下列步驟操作:

  1. 在抵達時確認 (DELETE 或 ACK) 目前的佇列訊息。在同一筆交易中,重新將新的佇列訊息排入佇列,並將遞送時間戳記設為未來的時間,超過處理所需的時間。
  2. 繼續處理,並在完成後確認新加入佇列的訊息。

這種做法的優點是不必持續延長租約,而且訊息不會重新傳送,直到未來時間到來為止 (涵蓋當機情況)。如果初始確認成功,即可實現最多一次處理。

為長時間執行的工作設定檢查點

Spanner 佇列可管理耗時數分鐘到數小時的工作,而不只是快速工作。如要執行這些長時間執行的工作,請使用下列方法:

  1. 在外部儲存中繼資料:使用頻外儲存空間來保存工作的詳細資料和狀態。
  2. 定期設定檢查點:如要從當機狀態復原,且不會遺失太多進度,工作應定期儲存自身狀態。
  3. 使用建議的檢查點模式:檢查點的最佳做法是自動確認 (ACK) 目前的佇列訊息,並傳送排定在未來傳送的新訊息。這則新訊息包含或指向更新後的狀態,可避免立即重新傳送至其他工作站。

即使無法進行完整檢查點作業,這種模式也能減少重複工作,但如果發生當機情況,工作會從頭開始。

排定在未來的特定時間執行工作

如要在未來特定時間安排工作,請在插入訊息時設定 DeliverTime 欄。

例如,試用期即將結束的提醒:

GoogleSQL

-- 1. Insert into Users table
INSERT INTO Users (UserId, UserName) VALUES (125, 'Trial User');

-- 2. Send message to queue with a future delivery time
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (125, 'trial-expire-reminder', b'{"type": "reminder"}', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 29 DAY));

PostgreSQL

-- 1. Insert into users table
INSERT INTO users (userid, username) VALUES (125, 'Trial User');

-- 2. Send message to queue with a future delivery time
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (125, 'trial-expire-reminder', CAST('{"type": "reminder"}' AS bytea), CURRENT_TIMESTAMP + INTERVAL '29 DAY');

處理大型訊息酬載

如果訊息酬載過大,請使用頻外儲存空間模式。將大型酬載儲存在獨立資料表中,並在佇列訊息中加入參照。

例如圖片處理:

GoogleSQL

-- Schema
CREATE TABLE ImageUploads (
  UserId    INT64 NOT NULL,
  ImageId   STRING(36) NOT NULL,
  ImageData BYTES(MAX),
  Status    STRING(MAX) -- PENDING, PROCESSING, DONE
) PRIMARY KEY (UserId, ImageId),
  INTERLEAVE IN PARENT Users;

CREATE QUEUE ImageProcessingQueue (
  UserId    INT64 NOT NULL,
  ImageId   STRING(36) NOT NULL,
  Payload   BYTES(1) NOT NULL -- Payload can be minimal
) PRIMARY KEY (UserId, ImageId),
  INTERLEAVE IN PARENT ImageUploads ON DELETE CASCADE;

-- Application Logic
-- 1. Upload image, insert into ImageUploads with Status 'PENDING'
-- 2. Send message to ImageProcessingQueue
INSERT INTO ImageProcessingQueue (UserId, ImageId, Payload) VALUES (123, 'image-uuid-1', b'');

-- Receiver for ImageProcessingQueue:
-- 1. Receives message (UserId, ImageId).
-- 2. Reads ImageData from ImageUploads.
-- 3. Processes image.
-- 4. Updates ImageUploads Status to 'DONE'.
-- 5. ACKs the queue message.

PostgreSQL

-- Schema
CREATE TABLE imageuploads (
  userid    bigint NOT NULL,
  imageid   varchar(36) NOT NULL,
  imagedata bytea,
  status    varchar, -- PENDING, PROCESSING, DONE
  PRIMARY KEY (userid, imageid)
) INTERLEAVE IN PARENT users;

CREATE QUEUE imageprocessingqueue (
  userid    bigint NOT NULL,
  imageid   varchar(36) NOT NULL,
  payload   bytea NOT NULL, -- Payload can be minimal
  PRIMARY KEY (userid, imageid)
) INTERLEAVE IN PARENT imageuploads ON DELETE CASCADE;

-- Application Logic
-- 1. Upload image, insert into imageuploads with status 'PENDING'
-- 2. Send message to imageprocessingqueue
INSERT INTO imageprocessingqueue (userid, imageid, payload) VALUES (123, 'image-uuid-1', CAST('' AS bytea));

-- Receiver for imageprocessingqueue:
-- 1. Receives message (userid, imageid).
-- 2. Reads imagedata from imageuploads.
-- 3. Processes image.
-- 4. Updates imageuploads status to 'DONE'.
-- 5. ACKs the queue message.

等待多個事件,再繼續操作

如要在繼續作業前等待多個事件 (例如聯結作業),請使用表格追蹤狀態,並使用佇列觸發檢查。

舉例來說,需要商品目錄和付款資訊的訂單履行作業:

  1. 建立含有 InventoryStatus 和 PaymentStatus 的 Orders 資料表。
  2. 確認庫存後,請更新 Orders,並傳送訊息到 OrderCheckQueue。
  3. 確認付款後,請更新 Orders 並傳送訊息到 OrderCheckQueue。
  4. OrderCheckQueue的收件者會檢查 Orders 表格。如果兩個狀態都已確認,系統會繼續出貨並確認訊息。如果沒有,系統可能會重新排隊,稍後再檢查,或執行其他邏輯。

定期執行動作

如要定期執行動作,請使用週期性排程模式。接收者確認訊息,並傳送排定在下一個間隔傳送的新訊息。

舉例來說,每小時資料匯總:

GoogleSQL

-- Inside your application transaction:
-- 1. Acknowledge current message
DELETE FROM AggregationQueue
WHERE TaskType = 'hourly-aggregator' AND MessageId = 'current-uuid'
ASSERT_ROWS_MODIFIED 1;

-- 2. Schedule next run 1 hour in the future
INSERT INTO AggregationQueue (TaskType, MessageId, Payload, DeliverTime)
VALUES ('hourly-aggregator', 'next-uuid', b'', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR));

PostgreSQL

-- Inside your application transaction:
-- 1. Acknowledge current message
DELETE FROM aggregationqueue
WHERE tasktype = 'hourly-aggregator' AND messageid = 'current-uuid'
ASSERT_ROWS_MODIFIED 1;

-- 2. Schedule next run 1 hour in the future
INSERT INTO aggregationqueue (tasktype, messageid, payload, deliver_time)
VALUES ('hourly-aggregator', 'next-uuid', CAST('' AS bytea), CURRENT_TIMESTAMP + INTERVAL '1 HOUR');

或者,您也可以使用用戶端程式庫的 Ack 和 Send 突變。這些範例假設您有一個封裝金鑰和酬載的 Message 物件:

Java

// Receiver logic for AggregationQueue
public void process(DatabaseClient dbClient, Message msg) {
  // ... do aggregation ...

  // ACK current message and schedule next run (1 hour from now)
  Instant nextRun = Instant.now().plus(Duration.ofHours(1));
  Mutation ackMutation =
      Mutation.newAckBuilder("AggregationQueue")
          .setKey(msg.getKey()) // Ack
          .build();
  Mutation sendMutation =
      Mutation.newSendBuilder("AggregationQueue")
          .setKey(Key.of("hourly-aggregator", "next-uuid"))
          .setPayload(Value.bytes(ByteArray.copyFrom("")))
          .setDeliveryTime(nextRun) // Schedule next
          .build();
  dbClient.write(Arrays.asList(ackMutation, sendMutation));
}

Go

// Receiver logic for AggregationQueue
func process(msg) {
    // ... do aggregation ...

    // ACK current message and schedule next run
    nextRun := time.Now().Add(1 * time.Hour)
    _, err := client.Apply(ctx, []*spanner.Mutation{
        spanner.Ack("AggregationQueue", msg.Key), // Ack
        spanner.Send("AggregationQueue",
            spanner.Key{"hourly-aggregator", "next-uuid"},
            []byte(""),
            spanner.WithDeliveryTime(nextRun), // Schedule next
        ),
    })
    // ... handle err ...
}

Python

# Receiver logic for AggregationQueue
def process(database: spanner.Database, msg: Message):
  # ... do aggregation ...
  # ACK current message and schedule next run (1 hour from now)
  next_run = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(
      hours=1
  )
  with database.batch() as batch:
    batch.ack(
        queue="AggregationQueue",
        key=msg.key,  # Ack
    )
    batch.send(
        queue="AggregationQueue",
        key=("hourly-aggregator", "next-uuid"),
        payload=b"",
        deliver_time=next_run,  # Schedule next
    )

Node.js

/**
 * Receiver logic for AggregationQueue
 * @param {import('@google-cloud/spanner').Database} database
 * @param { { key: Array<string|number>, payload: Buffer } } msg
 */
async function process(database, msg) {
  // ... do aggregation ...
  // ACK current message and schedule next run (1 hour from now)
  const nextRun = new Date(Date.now() + 60 * 60 * 1000);
  await database.runTransactionAsync(async (transaction) => {
    // Ack current message
    transaction.queueAck('AggregationQueue', msg.key);
    // Schedule next run
    transaction.queueSend(
      'AggregationQueue',
      ['hourly-aggregator', 'next-uuid'],
      {
        payload: Buffer.from(''),
        deliverTime: nextRun,
      }
    );
    await transaction.commit();
  });
}

使用暫時批次處理功能批次處理訊息

Spanner 佇列會盡可能縮短訊息傳送延遲時間。不過,如果獨立用戶端持續傳送大量訊息,個別處理每則訊息可能會造成高交易負荷。嘗試手動查詢或掃描佇列資料表來批次處理訊息,可能會導致範圍鎖定爭用、中止率提高,以及額外的讀取成本。

如要實現高處理量的批次處理,且不會發生爭用情形,請套用時間批次處理模式。傳送者會將訊息的 DeliverTime 對齊未來某個離散時間範圍 (例如,向上捨入至最接近的 10 秒界線)。由於獨立寄件者會計算相同的未來時間戳記,Spanner 傾向於將來自同一分割的訊息分組,並在max_batch_sizeRECEIVE_QUEUE_NAME()允許的情況下,以單一批次傳送這些訊息;如果 max_batch_size 小於要傳送的訊息數量,則會以多個批次傳送。

你可以使用下列公式計算傳送時間戳記:

$$ \text{DeliverTime} = \text{RoundDown}(\text{CurrentTime}, \text{FixedDelay}) + \text{FixedDelay} $$

舉例來說,如果時間範圍為 10 秒,在 09:05:00 到 09:05:09.999 之間加入佇列的訊息,都會收到 09:05:10 的 DeliverTime:

GoogleSQL

-- Calculate delivery time rounded to the next 10-second interval.
-- Use DIV(..., 10) * 10 to perform the RoundDown in SQL:
INSERT INTO OrderProcessingQueue (OrderId, DeliverTime, Payload)
VALUES (
  'order-101',
  TIMESTAMP_SECONDS(DIV(UNIX_SECONDS(CURRENT_TIMESTAMP()), 10) * 10 + 10),
  b'{"item": "book", "qty": 1}'
);

PostgreSQL

-- Calculate delivery time rounded to the next 10-second interval:
INSERT INTO orderprocessingqueue (orderid, deliver_time, payload)
VALUES (
  'order-101',
  to_timestamp((floor(extract(epoch from CURRENT_TIMESTAMP) / 10) * 10) + 10),
  CAST('{"item": "book", "qty": 1}' AS bytea)
);

接收者隨後會使用 max_batch_size 分批提取這些同步訊息:

GoogleSQL

SELECT OrderId, Payload, DeliverTime, SpannerLeaseToken
FROM RECEIVE_OrderProcessingQueue(max_duration=>'20m', max_batch_size=>50);

PostgreSQL

SELECT orderid, payload, deliver_time, spanner_lease_token
FROM spanner.receive_orderprocessingqueue(
    max_batch_size=>50, priority=>NULL, max_duration=>'20m');

如果同時傳送數百萬則訊息,將所有訊息對齊到完全相同的秒數,可能會導致處理量突然暴增。如要平均分配工作,並依實體批次處理訊息,請根據專屬 ID (例如 TIMESTAMP_ADD(deliver_time, INTERVAL MOD(ABS(FARM_FINGERPRINT(CAST(UserId AS STRING))), 10) SECOND)) 在計算中加入偏移量。

將資料修改與處理作業分離 (髒旗標模式)

在高輸送量的交易應用程式中,直接在面向使用者的交易中執行複雜的重新計算、搜尋索引或快取失效,可能會增加延遲,並導致共用資料列發生鎖定爭用,進而降低使用者體驗。

髒旗標模式可將資料修改作業與非同步處理作業分離。交易修改資料表時,會在同一交易的佇列中寫入輕量級的「髒位元」訊息。背景工作人員隨後會取用訊息,並以非同步方式執行耗費資源的處理作業。

舉例來說,當顧客變更個人資料或設定時:

GoogleSQL

-- Inside user profile update transaction:
-- 1. Update the primary entity table
UPDATE UserProfiles
SET FullName = 'Jane Doe', UpdatedAt = CURRENT_TIMESTAMP()
WHERE UserId = 456;

-- 2. Send lightweight dirty flag message to the queue.
-- It is recommended to interleave the queue in the primary UserProfiles
-- table for better transaction performance.
INSERT INTO UserDirtyQueue (UserId, TaskType, CommitTimestamp, Payload)
VALUES (456, 'reindex-user-profile', CURRENT_TIMESTAMP(), b'');

PostgreSQL

-- Inside user profile update transaction:
-- 1. Update the primary entity table
UPDATE userprofiles
SET fullname = 'Jane Doe', updatedat = CURRENT_TIMESTAMP
WHERE userid = 456;

-- 2. Send lightweight dirty flag message to the queue
INSERT INTO userdirtyqueue (userid, tasktype, committimestamp, payload)
VALUES (456, 'reindex-user-profile', CURRENT_TIMESTAMP, CAST('' AS bytea));

UserDirtyQueue 的背景接收器會接收 UserId,在使用者重要路徑之外讀取新的設定檔列,並重新計算搜尋索引或更新外部快取。如果短時間內將多個針對同一使用者的更新傳送至佇列,也可以將批次處理套用至這個模式。在這種情況下,接收端 TVF 可以指定大於 1 的 max_batch_size,從同一批次接收多則訊息。

監控工作站健康狀態並偵測逾時

您可以運用排定的佇列訊息,為工作節點或微服務執行個體機群建立容錯心跳和健康狀態監控系統。

如要實作健康狀態檢查,請按照下列步驟操作:

  1. 在啟動時註冊 worker:worker 初始化時,會將心跳訊息插入健康檢查佇列,並將未來的 DeliverTime 設為失敗期限 (例如 60 秒)。
  2. 定期傳送心跳訊號:在正常運作期間,Worker 會定期 (例如每 10 秒) 將 DeliverTime 往前推進 60 秒,藉此更新心跳訊息。
  3. 偵測失敗:如果工作人員當機或失去網路連線,心跳更新就會停止。60 秒後,傳送時間戳記會成熟 (DeliverTime <= CURRENT_TIMESTAMP()),Spanner 會將訊息傳送給警示接收者,啟動容錯移轉或工作重新指派。

重要事項:Spanner 佇列不支援 UPDATE DML 陳述式。因此,如要重新整理心跳時間戳記,您必須刪除現有訊息,並在單一交易中插入含有新 DeliverTime 的替代訊息,或是套用用戶端程式庫 Ack 和 Send 突變。請確保佇列的主鍵是 WorkerId (而非 (WorkerId, MessageId)),這樣每個 worker 在任何時間都只會有一則心跳訊息。

GoogleSQL

-- Inside the worker heartbeat transaction (executed every 10 seconds):
-- 1. Acknowledge the existing heartbeat message
DELETE FROM WorkerHealthQueue
WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

-- 2. Send replacement heartbeat with refreshed 60-second deadline
INSERT INTO WorkerHealthQueue (WorkerId, DeliverTime, Payload)
VALUES (
  'worker-node-42',
  TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 60 SECOND),
  b'{"status": "healthy", "active_jobs": 3}'
);

PostgreSQL

-- Inside the worker heartbeat transaction (executed every 10 seconds):
-- 1. Acknowledge the existing heartbeat message
DELETE FROM workerhealthqueue
WHERE workerid = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

-- 2. Send replacement heartbeat with refreshed 60-second deadline
INSERT INTO workerhealthqueue (workerid, deliver_time, payload)
VALUES (
  'worker-node-42',
  CURRENT_TIMESTAMP + INTERVAL '60 SECOND',
  CAST('{"status": "healthy", "active_jobs": 3}' AS bytea)
);

使用用戶端程式庫修改作業時,請確保修改作業切片中的 Ack 位於 Send 之前,如以下 Go 範例所示:

// Worker heartbeat loop
func sendHeartbeat(ctx context.Context, client *spanner.Client, workerID string) error {
    newDeadline := time.Now().Add(60 * time.Second)
    _, err := client.Apply(ctx, []*spanner.Mutation{
        spanner.Ack("WorkerHealthQueue", spanner.Key{workerID}),
        spanner.Send(
            "WorkerHealthQueue",
            spanner.Key{workerID},
            []byte(`{"status":"healthy"}`),
            spanner.WithDeliveryTime(newDeadline),
        ),
    })
    return err
}

工作站正常關機時,會明確刪除心跳訊息,因此不會觸發錯誤警報:

DELETE FROM WorkerHealthQueue WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

在單一佇列中處理多種工作類型

Spanner 執行個體的佇列總數有限制。為每個小型非同步作業建立個別佇列,很快就會達到這個限制,而且需要管理許多並行接收器查詢。

如要整合作業,請將不同工作類型合併至單一佇列,也就是多型佇列。建立多型佇列的策略有兩種。

策略 1:在主鍵中加入型別欄

GoogleSQL

CREATE QUEUE ApplicationTasks (
  TaskType   STRING(50) NOT NULL,
  TaskId     STRING(36) NOT NULL,
  Payload    BYTES(MAX) NOT NULL,
) PRIMARY KEY (TaskType, TaskId);

-- Enqueue an email task
INSERT INTO ApplicationTasks (TaskType, TaskId, Payload)
VALUES ('SEND_EMAIL', 'task-uuid-1', b'{"to": "user@example.com", "template": "welcome"}');

-- Enqueue an image thumbnail task
INSERT INTO ApplicationTasks (TaskType, TaskId, Payload)
VALUES ('GENERATE_THUMBNAIL', 'task-uuid-2', b'{"image_id": "img-789", "size": "small"}');

PostgreSQL

CREATE QUEUE applicationtasks (
  tasktype   varchar(50) NOT NULL,
  taskid     varchar(36) NOT NULL,
  payload    bytea NOT NULL,
  PRIMARY KEY (tasktype, taskid)
);

-- Enqueue an email task
INSERT INTO applicationtasks (tasktype, taskid, payload)
VALUES ('SEND_EMAIL', 'task-uuid-1', CAST('{"to": "user@example.com", "template": "welcome"}' AS bytea));

-- Enqueue an image thumbnail task
INSERT INTO applicationtasks (tasktype, taskid, payload)
VALUES ('GENERATE_THUMBNAIL', 'task-uuid-2', CAST('{"image_id": "img-789", "size": "small"}' AS bytea));

接收器會檢查 TaskType,並將酬載分派給對應的處理常式。

策略 2:多型酬載結構

或者,您也可以使用含有動作或型別鑑別子欄位的 JSON 酬載:

{
  "action": "SYNC_INVENTORY",
  "data": { "item_id": 987, "delta": -1 }
}

實作自訂重試延遲

Spanner 佇列會自動重試失敗或未確認的訊息,並內建指數輪詢機制。不過,如果訊息因已知原因而失敗,且失敗時間長度已知,或是因外部速率限制 (例如指定 Retry-After 標頭的 HTTP 429 回應) 而失敗,依賴自動退避可能會導致過早重試,浪費 CPU 資源。

如要實作自訂重試延遲時間,請按照下列步驟操作:

  1. 在訊息處理器中擷取特定暫時性失敗。
  2. 確認目前訊息,以滿足目前的傳送嘗試。
  3. 在同一筆交易中,傳送替代訊息,並將明確的 DeliverTime 設為所選的未來重試時間 (在下列範例中,為 5 分鐘後)。

GoogleSQL

-- Inside failure-handling transaction:
-- 1. Acknowledge the failed message
DELETE FROM OutboundNotificationQueue
WHERE NotificationId = 'notif-555' ASSERT_ROWS_MODIFIED 1;

-- 2. Reschedule delivery 5 minutes in the future
INSERT INTO OutboundNotificationQueue (NotificationId, DeliverTime, Payload)
VALUES (
  'notif-555',
  TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 5 MINUTE),
  b'{"recipient": "user@example.com", "retry_count": 2}'
);

PostgreSQL

-- Inside failure-handling transaction:
-- 1. Acknowledge the failed message
DELETE FROM outboundnotificationqueue
WHERE notificationid = 'notif-555' ASSERT_ROWS_MODIFIED 1;

-- 2. Reschedule delivery 5 minutes in the future
INSERT INTO outboundnotificationqueue (notificationid, deliver_time, payload)
VALUES (
  'notif-555',
  CURRENT_TIMESTAMP + INTERVAL '5 MINUTE',
  CAST('{"recipient": "user@example.com", "retry_count": 2}' AS bytea)
);

這種做法可讓應用程式精確管理退避時間表,並避免在下游復原期間耗盡外部 API 資源。

後續步驟