Skenario dan contoh antrean Spanner

Dokumen ini memberikan pola arsitektur dan contoh kode untuk skenario pengiriman pesan umum menggunakan antrean Spanner. Anda dapat menggunakan pola ini untuk memicu tugas asinkron setelah transaksi di-commit, menjadwalkan tugas yang tertunda atau berulang, mengelola payload pesan besar dengan penyimpanan di luar band, mengoordinasikan alur kerja multi-peristiwa, dan membuat checkpoint atau memperpanjang masa berlaku untuk tugas latar belakang yang berjalan lama.

Pemrosesan tepat satu kali dan pengakuan terima paling banyak satu kali

Berbagai pertimbangan dan solusi untuk pemrosesan tepat satu kali dan pengakuan paling banyak satu kali dijelaskan lebih mendetail di halaman Pemrosesan tepat satu kali dan pengakuan paling banyak satu kali.

Melakukan pekerjaan setelah transaksi di-commit

Untuk melakukan pekerjaan setelah transaksi di-commit, kirim pesan ke antrean dalam transaksi yang sama.

Misalnya, pendaftaran pengguna baru memicu email sambutan:

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)
);

Setelah transaksi dilakukan, penerima untuk UserTasks akan melakukan streaming pesan, mengirim email, dan mengonfirmasi pesan:

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(NULL, NULL, '20m');

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

Menangani pekerjaan yang berjalan lama

Jika Anda memiliki pekerjaan yang mungkin memerlukan waktu lebih lama dari masa sewa default (lebih dari 10 detik), panggil SELECT * FROM RENEWLEASE_QUEUE_NAME() secara berkala.

Misalnya, membuat laporan:

  1. Penerima akan mendapatkan pesan dari RECEIVE_ReportQueue().
  2. Mulai pembuatan laporan.
  3. Setiap 5 detik, panggil SELECT * FROM RENEWLEASE_ReportQueue([leaseToken]) dalam thread atau rutin terpisah.
  4. Setelah selesai, konfirmasi pesan dan simpan laporan.

Atau, jika Anda memiliki pekerjaan yang berjalan lama yang memerlukan pemrosesan paling banyak sekali, atau waktu sewa yang lama, lakukan hal berikut:

  1. Mengonfirmasi (DELETE atau ACK) pesan antrean saat ini saat tiba. Dalam transaksi yang sama, masukkan kembali pesan antrean baru dengan stempel waktu pengiriman di masa mendatang, di luar waktu yang diperlukan untuk pemrosesan.
  2. Lanjutkan pemrosesan, dan konfirmasi pesan yang baru dimasukkan ke dalam antrean setelah selesai.

Keuntungan dari pendekatan ini adalah tidak perlu terus-menerus memperpanjang masa sewa, dan pesan tidak dikirim ulang hingga waktu mendatang tiba (yang mencakup error). Jika pengakuan awal berhasil, maka akan tercapai pemrosesan paling banyak satu kali.

Membuat checkpoint pekerjaan yang berjalan lama

Antrean Spanner dapat mengelola tugas yang berlangsung selama beberapa menit hingga beberapa jam, bukan hanya tugas cepat. Untuk tugas yang berjalan lama ini, gunakan pendekatan berikut:

  1. Simpan metadata secara eksternal: Gunakan penyimpanan di luar band untuk menyimpan detail dan status tugas.
  2. Lakukan checkpoint secara rutin: Untuk memulihkan dari error tanpa kehilangan banyak progres, tugas harus menyimpan statusnya secara berkala.
  3. Gunakan pola pembuatan checkpoint yang direkomendasikan: Cara terbaik untuk membuat checkpoint adalah dengan mengonfirmasi (ACK) pesan antrean saat ini secara atomik dan mengirim pesan baru yang dijadwalkan untuk pengiriman pada masa mendatang. Pesan baru ini berisi atau mengarah ke status yang diperbarui, yang mencegah pengiriman ulang langsung ke worker lain.

Pola ini mengurangi pekerjaan duplikat meskipun checkpointing penuh tidak memungkinkan, meskipun tugas dimulai ulang dari awal setelah terjadi error dalam skenario tersebut.

Menjadwalkan tugas untuk waktu tertentu pada masa mendatang

Untuk menjadwalkan pekerjaan pada waktu tertentu di masa mendatang, tetapkan kolom DeliverTime saat menyisipkan pesan.

Misalnya, pengingat masa uji coba berakhir:

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');

Menangani payload pesan besar

Jika payload pesan Anda besar, gunakan pola penyimpanan di luar band. Simpan payload besar dalam tabel terpisah dan masukkan referensinya ke dalam pesan antrean.

Misalnya, pemrosesan gambar:

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.

Menunggu beberapa peristiwa sebelum melanjutkan

Untuk menunggu beberapa peristiwa sebelum melanjutkan (seperti operasi gabungan), gunakan tabel untuk melacak status dan antrean untuk memicu pemeriksaan.

Misalnya, pemenuhan pesanan yang memerlukan inventaris dan pembayaran:

  1. Buat tabel Orders dengan InventoryStatus dan PaymentStatus.
  2. Saat inventaris dikonfirmasi, perbarui Orders dan kirim pesan ke OrderCheckQueue.
  3. Setelah pembayaran dikonfirmasi, perbarui Orders dan kirim pesan ke OrderCheckQueue.
  4. Penerima untuk OrderCheckQueue memeriksa tabel Orders. Jika kedua status dikonfirmasi, penjual akan melanjutkan pengiriman dan mengonfirmasi pesan. Jika tidak, permintaan mungkin dimasukkan kembali ke antrean untuk pemeriksaan nanti atau menjalankan logika lain.

Melakukan tindakan secara berkala

Untuk melakukan tindakan secara berkala, gunakan pola penjadwalan berkala. Penerima mengonfirmasi pesan dan mengirim pesan baru yang dijadwalkan untuk interval berikutnya.

Misalnya, agregasi data per jam:

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');

Atau, gunakan mutasi Ack dan Send library klien. Contoh ini mengasumsikan Anda memiliki objek Message yang merangkum kunci dan payload:

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();
  });
}

Mengelompokkan pesan menggunakan pengelompokan temporal

Antrean Spanner mengirimkan pesan dengan latensi minimal. Namun, jika volume pesan yang tinggi terus-menerus datang dari klien independen, memproses setiap pesan satu per satu dapat menimbulkan overhead transaksi yang tinggi. Mencoba membuat kueri atau memindai tabel antrean secara manual untuk mengelompokkan pesan dapat menyebabkan pertentangan penguncian rentang, peningkatan rasio pembatalan, dan biaya baca tambahan.

Untuk mencapai batching throughput tinggi tanpa pertentangan, terapkan pola batching temporal. Pengirim menyelaraskan DeliverTime pesan ke periode waktu diskrit di masa mendatang (misalnya, membulatkan ke batas 10 detik terdekat). Karena pengirim independen menghitung stempel waktu mendatang yang identik, Spanner cenderung mengelompokkan pesan dari pemisahan yang sama dan mengirimkannya dalam satu batch jika max_batch_size dalam fungsi bernilai tabel RECEIVE_QUEUE_NAME() memungkinkan, atau beberapa batch jika max_batch_size lebih kecil dari jumlah pesan yang akan dikirim.

Stempel waktu pengiriman dapat dihitung dengan rumus ini:

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

Misalnya, dengan rentang waktu 10 detik, pesan yang dimasukkan dalam antrean antara 09.05.00 dan 09.05.09.999 semuanya menerima DeliverTime 09.05.10:

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)
);

Penerima kemudian menarik pesan yang disinkronkan ini dalam batch menggunakan 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');

Jika Anda mengirim jutaan pesan secara bersamaan, menyelaraskan semua pesan ke detik yang sama persis dapat menyebabkan lonjakan pemrosesan yang tiba-tiba. Untuk mendistribusikan pekerjaan secara merata dan tetap mengelompokkan pesan menurut entitas, tambahkan offset ke perhitungan berdasarkan ID unik (seperti TIMESTAMP_ADD(deliver_time, INTERVAL MOD(ABS(FARM_FINGERPRINT(CAST(UserId AS STRING))), 10) SECOND)).

Memisahkan modifikasi data dari pemrosesan (pola tanda kotor)

Dalam aplikasi transaksional dengan throughput tinggi, menjalankan penghitungan ulang yang kompleks, pengindeksan penelusuran, atau pembatalan validasi cache langsung di dalam transaksi yang terlihat oleh pengguna dapat memperlambat pengalaman pengguna dengan meningkatkan latensi dan menyebabkan persaingan penguncian pada baris bersama.

Pola tanda kotor memisahkan modifikasi data dari pemrosesan asinkron. Saat transaksi mengubah tabel, transaksi tersebut akan menulis pesan "bit kotor" ringan ke dalam antrean dalam transaksi yang sama. Kemudian, pekerja latar belakang menggunakan pesan dan melakukan pemrosesan yang berat secara asinkron.

Misalnya, saat pelanggan membuat perubahan profil atau setelan:

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));

Penerima latar belakang untuk UserDirtyQueue menerima UserId, membaca baris profil baru di luar jalur penting pengguna, dan menghitung ulang indeks penelusuran atau memperbarui cache eksternal. Pengelompokan dapat diterapkan pada pola ini juga jika beberapa pembaruan terhadap pengguna yang sama dikirim ke antrean dalam jangka waktu singkat. Dalam hal ini, TVF penerima dapat menentukan max_batch_size yang lebih besar dari 1 untuk menerima beberapa pesan dari batch yang sama.

Memantau kualitas pekerja dan mendeteksi waktu tunggu

Anda dapat menggunakan pesan antrean terjadwal untuk membangun sistem detak jantung dan pemantauan kondisi yang fault-tolerant untuk kumpulan node pekerja atau instance microservice.

Untuk menerapkan health check:

  1. Mendaftarkan pekerja saat startup: Saat diinisialisasi, pekerja akan memasukkan pesan detak jantung ke dalam antrean pemeriksaan kondisi dengan DeliverTime yang ditetapkan ke batas waktu kegagalannya (misalnya, 60 detik).
  2. Mengirim detak jantung berkala: Saat dalam kondisi baik, pekerja akan memperbarui pesan detak jantungnya secara berkala (misalnya, setiap 10 detik) dengan memajukan DeliverTime 60 detik ke depan.
  3. Mendeteksi kegagalan: Jika pekerja mengalami error atau kehilangan konektivitas jaringan, pembaruan detak jantung akan berhenti. Setelah 60 detik, stempel waktu pengiriman akan menjadi matang (DeliverTime <= CURRENT_TIMESTAMP()), dan Spanner akan mengirimkan pesan ke penerima pemberitahuan, yang memulai failover atau penugasan ulang tugas.

Penting: Antrean Spanner tidak mendukung pernyataan DML UPDATE. Oleh karena itu, untuk memperbarui stempel waktu detak jantung, Anda harus menghapus pesan yang ada dan menyisipkan pengganti dengan DeliverTime baru dalam satu transaksi, atau menerapkan mutasi library klien Ack dan Send. Pastikan kunci utama antrean hanya WorkerId (bukan (WorkerId, MessageId)) sehingga hanya ada satu pesan detak jantung per worker pada waktu tertentu.

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)
);

Saat menggunakan mutasi library klien, pastikan Ack mendahului Send dalam slice mutasi, seperti dalam contoh Go ini:

// 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
}

Saat worker dimatikan dengan benar, worker akan menghapus pesan detak jantungnya secara eksplisit sehingga tidak ada pemberitahuan palsu yang dipicu:

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

Menangani beberapa jenis tugas dalam satu antrean

Instance Spanner memiliki batas pada jumlah total antrean. Membuat antrean terpisah untuk setiap operasi asinkron kecil dapat dengan cepat mencapai batas ini dan memerlukan pengelolaan banyak kueri penerima serentak.

Untuk menggabungkan operasi, gabungkan berbagai jenis tugas ke dalam satu antrean, yang disebut antrean polimorfik. Ada dua strategi yang digunakan untuk membuat antrean polimorfik.

Strategi 1: Menyertakan kolom jenis dalam kunci utama

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));

Penerima memeriksa TaskType dan mengirimkan payload ke handler yang sesuai.

Strategi 2: Struktur payload polimorfik

Atau, gunakan payload JSON yang berisi kolom diskriminator tindakan atau jenis:

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

Menerapkan penundaan percobaan ulang kustom

Spanner Queues otomatis mencoba lagi pesan yang gagal atau tidak diakui dengan backoff eksponensial bawaan. Namun, dalam skenario saat pesan gagal karena alasan yang diketahui dengan durasi yang diketahui, atau batas frekuensi eksternal (seperti respons HTTP 429 yang menentukan Retry-After header), mengandalkan jeda otomatis dapat menyebabkan upaya percobaan ulang yang terlalu dini yang membuang-buang sumber daya CPU.

Untuk menerapkan penundaan percobaan ulang kustom:

  1. Tangkap kegagalan sementara tertentu di pemroses pesan Anda.
  2. Mengonfirmasi pesan saat ini untuk memenuhi upaya pengiriman saat ini.
  3. Dalam transaksi yang sama, kirim pesan pengganti dengan DeliverTime eksplisit yang ditetapkan ke waktu percobaan ulang di masa mendatang yang dipilih (dalam contoh berikut, 5 menit kemudian).

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)
);

Pendekatan ini memungkinkan aplikasi Anda mengelola jadwal mundur secara tepat dan menghindari saturasi API eksternal selama periode pemulihan hilir.

Langkah berikutnya