תרחישים ודוגמאות לתורים ב-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(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() באופן תקופתי.

לדוגמה, יצירת דוח:

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

טיפול בהודעות גדולות

אם מטען הייעודי (payload) של ההודעה גדול, כדאי להשתמש בתבנית אחסון מחוץ לפס. מאחסנים את המטען הייעודי (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.

המתנה למספר אירועים לפני שממשיכים

כדי להמתין למספר אירועים לפני שממשיכים (למשל, פעולת הצטרפות), משתמשים בטבלה למעקב אחר מצב ובמספר כדי להפעיל בדיקות.

לדוגמה, מילוי הזמנה שדורש מלאי ותשלום:

  1. תצור טבלה של Orders עם InventoryStatus ו-PaymentStatus.
  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 שמכיל את המפתח ואת המטען הייעודי (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();
  });
}

המאמרים הבאים