Menggunakan antrean Spanner

Dokumen ini menjelaskan cara menggunakan antrean Spanner. Dokumen ini menjelaskan cara membuat antrean, mengirim dan menerima pesan, memperpanjang masa berlaku pesan, dan mengonfirmasi pesan. Bagian ini juga mencakup praktik terbaik, informasi tentang pemantauan antrean, dan panduan pemecahan masalah.

Membuat antrean

Untuk membuat antrean, gunakan pernyataan CREATE QUEUE.

GoogleSQL

-- Example table for interleaving
CREATE TABLE Users (
  UserId   INT64 NOT NULL,
  UserName STRING(MAX)
) PRIMARY KEY (UserId);

-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE UserTasks (
  UserId     INT64 NOT NULL,
  MessageId  STRING(36) NOT NULL, -- UUID recommended
  Payload    BYTES(MAX) NOT NULL  -- Also: Proto, JSON, String are possible.
) PRIMARY KEY (UserId, MessageId),
INTERLEAVE IN PARENT Users ON DELETE CASCADE;

PostgreSQL

-- Example table for interleaving
CREATE TABLE users (
  userid   bigint NOT NULL,
  username varchar,
  PRIMARY KEY (userid)
);

-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE usertasks (
  userid     bigint NOT NULL,
  messageid  varchar(36) NOT NULL, -- UUID recommended
  payload    bytea NOT NULL, -- Also: text, varchar, jsonb are possible.
  PRIMARY KEY (userid, messageid)
) INTERLEAVE IN PARENT users ON DELETE CASCADE;

Satu-satunya kolom yang tidak perlu dibuat secara eksplisit dalam pernyataan CREATE QUEUE disebut DeliverTime di GoogleSQL dan deliver_time di PostgreSQL. Indeks ini dibuat secara otomatis oleh Spanner.

Antrean mendukung kebijakan time to live (TTL), yang dapat membantu mengelola backlog pesan untuk pesan lama yang belum dikonfirmasi.

Kirim pesan

Untuk mengirim pesan ke antrean, gunakan pernyataan INSERT DML:

GoogleSQL

-- Send message immediately
INSERT INTO Users (UserId) VALUES (123);
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-1', b'Your task payload here', CURRENT_TIMESTAMP());

-- Schedule a message delivery
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-2', b'Scheduled task', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR));

PostgreSQL

-- Send message immediately
INSERT INTO users (userid) VALUES (123);
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-1', CAST('Your task payload here' AS bytea), CURRENT_TIMESTAMP);

-- Schedule a message delivery
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-2', CAST('Scheduled task' AS bytea), CURRENT_TIMESTAMP + INTERVAL '1 HOUR');

Atau, gunakan mutasi Send library klien untuk menyisipkan pesan:

Go

m := spanner.Send("UserTasks", spanner.Key{int64(123), "some-unique-id-1"}, []byte("Your task payload here"), spanner.WithDeliveryTime(futureTime))
_, err := client.Apply(ctx, []*spanner.Mutation{m})

Java

dbClient.write(
        Collections.singletonList(
            Mutation.newSendBuilder("UserTasks")
                .setKey(Key.of(123L, "some-unique-id-1"))
                .setPayload(Value.bytes(ByteArray.copyFrom("message3")))
                .setDeliveryTime(futureTime)
                .build()));

Menerima pesan

Gunakan fungsi bernilai tabel (TVF) RECEIVE_QUEUE_NAME() dengan ExecuteStreamingSQL untuk menerima pesan. Ini adalah panggilan yang berjalan lama. Anda harus menjalankan salah satu panggilan ini per pekerja, per antrean, secara berulang. Jalankan kueri menggunakan pembacaan kuat karena Spanner menolak pembacaan yang sudah tidak berlaku.

GoogleSQL

-- SQL query to stream messages
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');

PostgreSQL

-- SQL query to stream messages
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token
FROM spanner.receive_usertasks(NULL, NULL, '20m');

Kode klien Anda harus melakukan iterasi pada hasil menggunakan kueri streaming. Setiap baris yang ditampilkan adalah pesan.

Menerima pesan dalam batch

Untuk meningkatkan throughput dengan memproses beberapa pesan secara bersamaan, Anda dapat menerima pesan dalam batch dengan menentukan argumen max_batch_size:

GoogleSQL

-- SQL query to stream messages
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken,
  SpannerLastBatchMessage -- Special boolean column returns TRUE if the
                          -- last message is in a batch.
FROM RECEIVE_UserTasks(max_duration=>'20m', max_batch_size=>20);

PostgreSQL

-- SQL query to stream messages
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token,
  spanner_last_batch_message -- Special boolean column returns TRUE if the
                             -- last message is in a batch.
FROM spanner.receive_usertasks(20, NULL, '20m');

Atau, gunakan library klien. Contoh Go ini menunjukkan cara melakukan streaming pesan dari antrean, memverifikasi masa berlaku sewa, dan mengonfirmasi pesan secara asinkron:

// import "cloud.google.com/go/spanner"
// import "google.golang.org/api/iterator"

stmt := spanner.Statement{SQL: "SELECT * FROM RECEIVE_UserTasks(max_duration=>'20m')"}
iter := client.Single().Query(ctx, stmt)
defer iter.Stop()

for {
    row, err := iter.Next()
    if err == iterator.Done {
        break // Or potentially restart the query
    }
    if err != nil {
        // Handle error
        return err
    }

    var userId int64
    var messageId string
    var payload []byte
    var deliverTime time.Time
    var leaseExpiration time.Time
    var leaseToken string
    var lastBatchMessage bool

    if err := row.Columns(&userId, &messageId, &payload, &deliverTime, &leaseExpiration, &leaseToken, &lastBatchMessage); err != nil {
        // Handle column parsing error
        return err
    }

    if time.Now().After(leaseExpiration) {
        log.Printf("Lease expired for message %s, skipping", messageId)
        continue
    }

    // Process and acknowledge the message asynchronously so that we don't
    // block receiving subsequent messages.
    go func(userId int64, messageId string, payload []byte) {
        // Process the message (payload)
        // ... potentially long-running work ...
        // Need to extend the lease if processing is long

        // Acknowledge the message upon success
        _, err := client.Apply(ctx, []*spanner.Mutation{
            spanner.Ack("UserTasks", spanner.Key{userId, messageId}),
        })
        if err != nil {
            // Handle ack error
        }
    }(userId, messageId, payload)
}

Memperpanjang masa sewa pesan

Jika pemrosesan pesan memerlukan waktu lebih lama daripada masa berlaku awal (10 detik), gunakan sintaksis berikut:

GoogleSQL

-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => ['token1', ..., 'tokenN'])

-- Returns rows of tokens and whether they were successfully extended
SpannerOldLeaseToken   SpannerNewLeaseToken  SpannerLeaseExpirationTimestamp
        <old_token1>           <new_token1>  "2025-09-27T12:10:00.0Z"
            'token2'                 <NULL>  "2025-09-27T12:09:51.0Z"
...
            'tokenN'           'tokenN_new'  "2025-09-27T12:10:01.0Z"

PostgreSQL

-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM spanner.renewlease_usertasks(lease_tokens => ARRAY['token1', ..., 'tokenN'])

-- Returns rows of tokens and whether they were successfully extended
spanner_old_lease_token   spanner_new_lease_token  spanner_lease_expiration_timestamp
           <old_token1>              <new_token1>  "2025-09-27T12:10:00.0Z"
               'token2'                    <NULL>  "2025-09-27T12:09:51.0Z"
...
               'tokenN'              'tokenN_new'  "2025-09-27T12:10:01.0Z"

Token sewa ditampilkan sesuai dengan logika berikut:

  1. Token sewa yang tidak dapat diuraikan tidak menampilkan baris.
  2. Token sewa yang sudah habis masa berlakunya tidak menampilkan baris.
  3. Token sewa yang tidak dapat diperpanjang akan menampilkan baris dengan SpannerNewLeaseToken NULL. Hal ini dapat terjadi jika pesan sudah dikonfirmasi, tetapi token sewa belum habis masa berlakunya.

Atau, gunakan library klien. Contoh Go ini menunjukkan cara memperpanjang masa berlaku pesan:

// Inside message processing loop...

// Before leaseExpiration, for example, in a separate goroutine or timed check
// leaseToken is from the SELECT query
extendStmt := spanner.Statement{
    SQL: "SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => [@token1])",
    Params: map[string]interface{}{
        "token1": leaseToken,
    },
}
_, err := client.Single().Query(ctx, extendStmt).Next() // Simplified call
if err != nil {
    log.Printf("Failed to extend lease for %s: %v", messageId, err)
    // Processing should probably stop as redelivery is likely
} else {
    // New lease expiration is typically approximately 10s from now
    log.Printf("Lease extended for %s", messageId)
    // Update local leaseExpiration time if needed
}

Mengonfirmasi pesan

Gunakan DML DELETE untuk mengonfirmasi pesan. Ini harus bersifat transaksional dengan penulisan lain yang terkait dengan pemrosesan pesan.

GoogleSQL

-- DML for acknowledging
DELETE FROM UserTasks WHERE UserId = @userId AND MessageId = @messageId ASSERT_ROWS_MODIFIED 1;

PostgreSQL

-- DML for acknowledging
DELETE FROM usertasks WHERE userid = $1 AND messageid = $2 ASSERT_ROWS_MODIFIED 1;

Atau, gunakan mutasi Ack library klien untuk mengonfirmasi pesan:

Go

_, err := client.Apply(ctx, []*spanner.Mutation{
    spanner.Ack("UserTasks", spanner.Key{1}),
})

Java

dbClient.write(
    Collections.singletonList(
        Mutation.newAckBuilder("UserTasks")
            .setKey(Key.of(2L))
            .build()));

Praktik terbaik

Berikut adalah praktik terbaik untuk menggunakan antrean Spanner:

  • Payload kecil: jaga agar payload pesan antrean tetap kecil di bawah 4 KB. Gunakan pola penyimpanan di luar band untuk data yang lebih besar.
  • Pengelolaan sewa: memperpanjang sewa untuk tugas yang mungkin melampaui sewa pesan default 10 detik. Kegagalan memperpanjang masa berlaku pesan dapat menyebabkan pengiriman ulang dan potensi pemrosesan ganda.
  • Penanganan error: Spanner mengantrekan pesan yang gagal diproses dengan penundaan dalam satu jam pertama, dan pesan yang lebih lama akan dicoba lagi sekali per jam. Pertimbangkan untuk memindahkan pesan yang gagal secara permanen ke antrean terpisah.
  • Pemantauan: pantau kedalaman antrean dan usia pesan terlama yang belum dikonfirmasi untuk memastikan penerima tidak membatasi penyerapan pipeline Anda.
  • Idempotensi: desain pemroses pesan Anda agar dapat mentoleransi operasi berulang, karena pengiriman minimal satu kali berarti duplikat sesekali mungkin terjadi. Lihat halaman Pemrosesan tepat satu kali dan konfirmasi paling banyak satu kali untuk mengetahui informasi selengkapnya.
  • Durasi fungsi bernilai tabel: hindari durasi yang sangat singkat dan terlalu lama untuk fungsi bernilai tabel. Durasi sedang, seperti 20 menit, direkomendasikan.
  • Ukuran batch: sesuaikan max_batch_size dengan workload Anda. Gunakan batch yang lebih kecil untuk peristiwa fan-out tinggi guna menghindari pertentangan kunci pada baris bersama. Gunakan batch yang lebih besar untuk tugas independen dengan kueri per detik yang tinggi. Mengonfirmasi atau memperpanjang masa berlaku pesan untuk batch dalam satu transaksi untuk performa terbaik.

Memantau

Anda dapat memantau operasi antrean menggunakan tabel introspeksi Spanner. Meskipun tabel ini tidak menyertakan kolom khusus antrean, Anda dapat mengidentifikasi aktivitas antrean dengan menelusuri nama antrean yang ditentukan pengguna di tabel berikut:

Metrik antrean Spanner berikut dapat ditemukan di bawah awalan spanner.googleapis.com/queue/* di Cloud Monitoring:

  • buffered_ready_messages: (GAUGE, INT64, 1) Jumlah pesan yang disimpan dalam memori dan siap dikirim ke penerima.
  • message_send_count: (DELTA, INT64, 1) Jumlah pesan yang dikirim di Spanner selama interval untuk antrean.
  • message_ack_count: (DELTA, INT64, 1) Jumlah pesan yang dikonfirmasi di Spanner selama interval untuk antrean.
  • oldest_unacked_message_age: (GAUGE, INT64, 1) Usia (dalam detik) pesan tertua yang belum dikonfirmasi dalam antrean.
  • lease_expiration_count: (DELTA, INT64, 1) Jumlah masa berlaku sewa yang berakhir di Spanner selama interval untuk antrean.

Semua metrik sebelumnya diambil sampelnya kira-kira setiap 60 detik. Setelah pengambilan sampel, data mungkin tidak terlihat hingga 120 detik. Pencatatan audit Spanner mencakup operasi tulis, baca, dan skema pada antrean.

Memecahkan masalah

Bagian berikut menjelaskan cara mengidentifikasi dan menyelesaikan masalah umum saat menggunakan antrean Spanner.

Backlog pesan bertambah

Diagnosis

oldest_unacked_message_age dan buffered_ready_messages ditinggikan. Hal ini menunjukkan ketidakseimbangan antara kecepatan pengiriman pesan dan kapasitas pemrosesan pesan aplikasi Anda.

Resolusi

Untuk mengatasi masalah ini, lakukan langkah berikut:

  • Periksa rasio pengakuan: jika metrik message_ack_count telah menurun, periksa pekerja klien Anda untuk memastikan bahwa mereka berjalan dengan benar dan tidak terhenti atau error.
  • Periksa kecepatan pengiriman: jika message_send_count meningkat tajam, tingkatkan skala worker pemrosesan pesan untuk menangani peningkatan beban.
  • Konfirmasi kehabisan resource: periksa apakah jumlah buffered_ready_messages dan lease_expiration_count meningkat. Kombinasi ini menunjukkan kekurangan penerima fungsi bernilai tabel (TVF) aktif atau pemrosesan klien yang lambat.

Pesan individual macet

Diagnosis

Metrik oldest_unacked_message_age tinggi, tetapi buffered_ready_messages rendah atau stabil. Hal ini menunjukkan bahwa pesan individual gagal diproses atau dikonfirmasi, bukan karena hambatan kapasitas secara keseluruhan.

Resolusi

Untuk mengatasi masalah ini, lakukan langkah berikut:

  • Mengidentifikasi pesan yang macet: kueri tabel antrean dan urutkan hasil menurut waktu pengiriman untuk menemukan pesan terlama yang belum dikonfirmasi:

    GoogleSQL

    SELECT *
    FROM UserTasks
    ORDER BY DeliverTime ASC
    LIMIT 10;
    

    PostgreSQL

    SELECT *
    FROM usertasks
    ORDER BY deliver_time ASC
    LIMIT 10;
    
  • Selidiki kegagalan pemrosesan: periksa log aplikasi Anda untuk menentukan alasan pekerja tidak mengonfirmasi pesan yang diidentifikasi. Pesan yang belum dikonfirmasi akan dikirim ulang secara otomatis setelah masa berlakunya berakhir.

Masa berlaku sewa pesan berakhir sebelum pemrosesan selesai

Diagnosis

Metrik lease_expiration_count meningkat atau naik. Hal ini menunjukkan bahwa waktu pemrosesan pesan melebihi durasi lease (default 10 detik) sebelum pekerja dapat mengonfirmasi pesan.

Resolusi

Untuk mengatasi masalah ini, lakukan langkah berikut:

  • Perpanjang masa berlaku sewa secara proaktif: jika pemrosesan pesan memerlukan waktu lebih dari 10 detik, panggil RENEWLEASE_QUEUE_NAME() TVF secara berkala. Perpanjang masa sewa sekitar 7-8 detik ke dalam pemrosesan untuk memperhitungkan latensi jaringan dengan aman.
  • Selidiki pemrosesan yang lambat: jika aplikasi Anda sudah memperpanjang masa berlaku sewa secara aktif, tetapi lease_expiration_count tetap tinggi, periksa kode backend Anda untuk menemukan hambatan pemrosesan, panggilan RPC yang lambat, atau kebuntuan.

Tidak dapat menskalakan kecepatan pengiriman pesan

Diagnosis

Metrik message_send_count mencapai dataran tinggi throughput atau permintaan publikasi mengalami peningkatan latensi penulisan saat Anda mencoba meningkatkan kecepatan pengiriman.

Resolusi

Untuk mengatasi masalah ini, pertimbangkan perubahan struktural berikut:

  • Meningkatkan skala resource komputasi: tambahkan node atau unit pemrosesan ke instance Spanner Anda untuk meningkatkan kapasitas database secara keseluruhan.
  • Periksa pemisahan: evaluasi apakah Anda dapat menambahkan lebih banyak pemisahan untuk mendistribusikan beban tulis di beberapa server.
  • Mengoptimalkan perutean: pastikan aplikasi penerbitan Anda menulis langsung ke region leader instance Spanner Anda untuk meminimalkan latensi penulisan.

Pesan yang siap tidak diproses

Diagnosis

Metrik buffered_ready_messages tinggi dan meningkat. Hal ini menunjukkan bahwa pesan di-buffer dan siap dikirim dalam memori, tetapi pekerja penerima tidak mengambilnya.

Resolusi

Untuk mengatasi masalah ini, lakukan langkah berikut:

  • Periksa koneksi TVF aktif: periksa jumlah TVF penerima aktif Anda untuk memastikan pekerja pembaca terhubung dan menarik pesan secara aktif. Jika pekerja terputus atau tidak menjalankan kueri serentakRECEIVE_QUEUE_NAME() yang cukup, pesan akan tetap tidak di-polling di buffer.

Langkah berikutnya