במסמך הזה מפורטים דפוסי ארכיטקטורה ודוגמאות קוד לתרחישים נפוצים של העברת הודעות באמצעות תורים ב-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(NULL, NULL, '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() באופן תקופתי.
לדוגמה, יצירת דוח:
- המקבל מקבל הודעה מ
RECEIVE_ReportQueue(). - מתחילים ליצור את הדוח.
- כל 5 שניות, קוראים ל-
SELECT * FROM RENEWLEASE_ReportQueue([leaseToken])בשרשור או בשגרה נפרדים. - אחרי שהדוח מוכן, מאשרים את ההודעה ושומרים את הדוח.
לחלופין, אם יש לכם עבודה שדורשת עיבוד לכל היותר פעם אחת או זמן השכרה ארוך, אתם יכולים לבצע את הפעולות הבאות:
- לאשר (
DELETEאוACK) את ההודעה הנוכחית בתור כשהיא מגיעה. באותה עסקה, מוסיפים מחדש לרשימת ההמתנה הודעה חדשה עם חותמת זמן למסירה בעתיד, מעבר לזמן שנדרש לעיבוד. - להמשיך בעיבוד ולאשר את ההודעה החדשה שהוכנסה לתור כשהעיבוד מסתיים.
היתרונות של הגישה הזו הם שאין צורך להאריך את תקופת ההשכרה באופן רציף, וההודעה לא נמסרת מחדש עד שהזמן העתידי מגיע (מה שמכסה קריסות). אם האישור הראשוני יצליח, העיבוד יתבצע לכל היותר פעם אחת.
נקודת ביקורת של עבודה ממושכת
תורים ב-Spanner יכולים לנהל משימות שנמשכות דקות עד שעות, ולא רק משימות מהירות. כדי לבצע משימות ארוכות, כדאי לפעול לפי הגישה הבאה:
- אחסון מטא-נתונים חיצוני: שימוש באחסון מחוץ לפס כדי לשמור את הפרטים והמצב של המשימה.
- נקודת ביקורת באופן קבוע: כדי לשחזר ממצבי קריסה בלי לאבד הרבה התקדמות, המשימה צריכה לשמור את המצב שלה באופן תקופתי.
- שימוש בתבנית המומלצת של נקודות ביקורת: הדרך הכי טובה להשתמש בנקודות ביקורת היא לאשר באופן אטומי (
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');
טיפול במטענים גדולים של הודעות
אם מטען הייעודי (payload) של ההודעה גדול, כדאי להשתמש בדפוס אחסון מחוץ לפס. מאחסנים את המטען הגדול בטבלה נפרדת ומכניסים הפניה אליו בהודעת התור.
לדוגמה, עיבוד תמונה:
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.
המתנה למספר אירועים לפני שממשיכים
כדי להמתין לכמה אירועים לפני שממשיכים (למשל, פעולת הצטרפות), משתמשים בטבלה כדי לעקוב אחרי המצב ובתור כדי להפעיל בדיקות.
לדוגמה, הכנת מוצר לאספקה שדורשת מלאי ותשלום:
- תצור טבלה של
OrdersעםInventoryStatusו-PaymentStatus. - אחרי אישור המלאי, מעדכנים את
Ordersושולחים הודעה לכתובתOrderCheckQueue. - אחרי אישור התשלום, צריך לעדכן את
Ordersולשלוח הודעה לכתובתOrderCheckQueue. - הנמען של
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 שמכיל את המפתח ואת המטען הייעודי (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();
});
}
שליחת הודעות בקבוצות באמצעות שליחת קבוצות זמניות
תורים ב-Spanner מעבירים הודעות עם זמן אחזור מינימלי. עם זאת, כשכמויות גדולות של הודעות מגיעות באופן רציף מלקוחות עצמאיים, עיבוד כל הודעה בנפרד עלול ליצור תקורה גבוהה של עסקאות. ניסיון לבצע שאילתה או סריקה של טבלת התור באופן ידני כדי לבצע עיבוד אצווה של הודעות עלול לגרום לתחרות על נעילת טווח, לשיעורי ביטול גבוהים יותר ולעלויות קריאה נוספות.
כדי להשיג עיבוד באצווה עם תפוקה גבוהה בלי התנגשות, צריך להשתמש בדפוס temporal batching. השולחים מתאימים את DeliverTime של ההודעות לחלון זמן נפרד בעתיד (לדוגמה, עיגול כלפי מעלה לגבול הקרוב ביותר של 10 שניות). מכיוון ששולחים עצמאיים מחשבים חותמת זמן עתידית זהה, מערכת Spanner נוטה לקבץ את ההודעות מאותו פיצול יחד ולשלוח אותן באצווה אחת אם הפונקציה max_batch_size בטבלה RECEIVE_QUEUE_NAME() מאפשרת זאת, או בכמה אצוות אם max_batch_size קטן ממספר ההודעות שצריך לשלוח.
אפשר לחשב את חותמת הזמן של המסירה באמצעות הנוסחה הבאה:
לדוגמה, אם חלון הזמן הוא 10 שניות, כל ההודעות שנוספו לתור בין השעות 09:05:00 ל-09:05:09.999 יקבלו את הערך 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)
);
אחר כך הנמענים מושכים את ההודעות האלה שמתזמנות את עצמן בקבוצות באמצעות 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');
אם שולחים מיליוני הודעות בו-זמנית, התאמה של כל ההודעות לאותה שנייה בדיוק עלולה לגרום לעליות פתאומיות בעיבוד. כדי לחלק את העבודה באופן שווה ועדיין לאגד הודעות לפי ישות, מוסיפים להיסט לחישוב על סמך מזהה ייחודי (כמו 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 של המקבל יכול להיות גדול מ-1max_batch_size כדי לקבל כמה הודעות מאותו אצווה.
מעקב אחרי תקינות העובדים וזיהוי פסק זמן
אתם יכולים להשתמש בהודעות בתור שמוגדרות למועד מסוים כדי ליצור מערכת עמידה בפני תקלות למעקב אחרי פעילות תקינה של קבוצות של צמתי עובד או מופעים של מיקרו-שירותים.
כדי להטמיע בדיקות תקינות:
- Register worker on startup: When a worker initializes, it inserts a
heartbeat message into a health-check queue with a future
DeliverTimeset to its failure deadline (for example, 60 seconds). - שליחת אותות מחזוריות: כשהעובד תקין, הוא מרענן את אות המחזור שלו באופן מחזורי (לדוגמה, כל 10 שניות) על ידי קידום
DeliverTimeבעוד 60 שניות קדימה. - זיהוי כשלים: אם תהליך העובד קורס או שהחיבור לרשת מתנתק, הרענון של אותות החיים מפסיק. אחרי 60 שניות, חותמת הזמן של המסירה מתעדכנת (
DeliverTime <= CURRENT_TIMESTAMP()), ו-Spanner מעביר את ההודעה לנמען שמוגדר לקבל התראות, שמתחיל בהעברה או בהקצאה מחדש של המשימה.
חשוב: תורי Spanner לא תומכים בהצהרות UPDATE DML. לכן, כדי לרענן את חותמת הזמן של אות החיים, צריך למחוק את ההודעה הקיימת ולהוסיף הודעה חלופית עם DeliverTime חדש בתוך עסקה אחת, או להחיל שינויים בספריות הלקוח Ack ו-Send. מוודאים שהמפתח הראשי של התור הוא WorkerId בלבד (ולא (WorkerId, MessageId)), כדי שתהיה רק הודעת דופק אחת לכל עובד בכל זמן נתון.
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 ושולח את המטען הייעודי (payload) ל-handler המתאים.
אסטרטגיה 2: מבנה מטען ייעודי משתנה
אפשרות אחרת היא להשתמש ב-payload של JSON שמכיל שדה של פעולה או של סוג מפלה:
{
"action": "SYNC_INVENTORY",
"data": { "item_id": 987, "delta": -1 }
}
הטמעה של השהיות מותאמות אישית של ניסיונות חוזרים
בתורים של Spanner, המערכת מנסה באופן אוטומטי לשלוח מחדש הודעות שנכשלו או שלא אושרו, עם השהיה מעריכית לפני ניסיון חוזר (exponential backoff) מובנית. עם זאת, בתרחישים שבהם הודעה נכשלת בגלל סיבה ידועה עם משך ידוע, או בגלל הגבלת קצב חיצונית (כמו תגובת HTTP 429 שמציינת כותרת Retry-After), הסתמכות על השהיה אוטומטית עלולה לגרום לניסיונות חוזרים מוקדמים מדי שמבזבזים משאבי CPU.
כדי להגדיר השהיה מותאמת אישית בין ניסיונות חוזרים:
- לזהות את הכשל הזמני הספציפי במעבד ההודעות.
- צריך לאשר את ההודעה הנוכחית כדי שהניסיון הנוכחי למסירה יצליח.
- באותה טרנזקציה, שולחים הודעת החלפה עם
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 חיצוניים במהלך תקופות ההתאוששות של הנתונים.
המאמרים הבאים
- איך משתמשים בתורים של Spanner, כולל שיטות מומלצות ומעקב
- מידע על עיבוד בדיוק פעם אחת ואישור לכל היותר פעם אחת
- הגדרת בקרת גישה באמצעות בקרת גישה פרטנית לתורים.