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:
- Token sewa yang tidak dapat diuraikan tidak menampilkan baris.
- Token sewa yang sudah habis masa berlakunya tidak menampilkan baris.
- Token sewa yang tidak dapat diperpanjang akan menampilkan baris dengan
SpannerNewLeaseTokenNULL. 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_sizedengan 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:
- Statistik operasi baca
- Statistik transaksi
- Statistik kunci
- Statistik ukuran tabel
- Statistik operasi tabel
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_counttelah menurun, periksa pekerja klien Anda untuk memastikan bahwa mereka berjalan dengan benar dan tidak terhenti atau error. - Periksa kecepatan pengiriman: jika
message_send_countmeningkat tajam, tingkatkan skala worker pemrosesan pesan untuk menangani peningkatan beban. - Konfirmasi kehabisan resource: periksa apakah jumlah
buffered_ready_messagesdanlease_expiration_countmeningkat. 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_counttetap 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 serentak
RECEIVE_QUEUE_NAME()yang cukup, pesan akan tetap tidak di-polling di buffer.
Langkah berikutnya
- Jelajahi lainnya skenario dan contoh antrean Spanner.
- Pelajari pemrosesan tepat satu kali dan konfirmasi paling banyak satu kali.
- Konfigurasi kontrol akses dengan kontrol akses terperinci untuk antrean.