שימוש בתורים של Spanner

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

יצירת רשימת 'הבאים בתור'

כדי ליצור תור, משתמשים בהצהרה 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;

העמודה היחידה שלא צריך ליצור באופן מפורש בהצהרת CREATE QUEUE נקראת DeliverTime ב-GoogleSQL ו-deliver_time ב-PostgreSQL. הן נוצרות באופן אוטומטי על ידי Spanner.

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

שליחת הודעה

כדי לשלוח הודעה לתור, משתמשים בהצהרת 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');

אפשר גם להשתמש במוטציות של ספריית הלקוח Send כדי להוסיף הודעה:

המשך

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

קבלת הודעות

משתמשים בפונקציה RECEIVE_QUEUE_NAME() שמחזירה טבלה (TVF) עם ExecuteStreamingSQL כדי לקבל הודעות. זו שיחה ארוכה. צריך להפעיל אחת מהשיחות האלה לכל עובד, לכל תור, בלולאה. מריצים את השאילתה באמצעות קריאה חזקה כי Spanner דוחה קריאות לא עדכניות.

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

קוד הלקוח צריך לחזור על התוצאות באמצעות שאילתת סטרימינג. כל שורה שמוחזרת היא הודעה.

קבלת הודעות באצווה

כדי להגדיל את קצב העברת הנתונים על ידי עיבוד של כמה הודעות יחד, אפשר לקבל הודעות באצוות על ידי ציון הארגומנט 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');

אפשרות אחרת היא להשתמש בספריית הלקוח. בדוגמה הבאה בשפת Go מוצגות הפעולות הבאות: הזרמת הודעות מתור, אימות של תפוגת ההרשאה ואישור הודעות באופן אסינכרוני:

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

הארכת תקופת ההשאלה של ההודעה

אם העיבוד של הודעה נמשך יותר מהזמן הראשוני להקצאת משאבים (10 שניות), משתמשים בתחביר הבא:

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"

אסימוני חכירה מוחזרים בהתאם ללוגיקה הבאה:

  1. טוקנים של חכירה שלא ניתן לנתח לא מחזירים שורה.
  2. אם תוקף האסימונים להשכרה כבר פג, לא תוחזר שורה.
  3. אסימוני ליסינג שלא ניתן לחדש מחזירים שורה עם SpannerNewLeaseToken של NULL. מצב כזה יכול לקרות אם ההודעה כבר אושרה אבל טוקן ההרשאה לא פג.

אפשרות אחרת היא להשתמש בספריית הלקוח. בדוגמה הזו בשפת Go אפשר לראות איך להאריך את תקופת ההחזקה של הודעה:

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

אישור קבלת הודעה

משתמשים ב-DELETE DML כדי לאשר קבלת הודעה. הפעולה הזו חייבת להיות טרנזקציונלית עם פעולות כתיבה אחרות שקשורות לעיבוד ההודעה.

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;

אפשרות אחרת היא להשתמש במוטציה Ack בספריית הלקוח כדי לאשר קבלת הודעה:

המשך

_, 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()));

שיטות מומלצות

השיטות המומלצות הבאות יעזרו לכם להשתמש בתורים של Spanner:

  • מטענים קטנים: מומלץ לשמור על מטענים קטנים של הודעות בתור, עד 4KB. כדי לאחסן נתונים גדולים יותר, כדאי להשתמש בתבנית אחסון מחוץ לפס.
  • ניהול פרק זמן לעיבוד: הארכת פרק זמן לעיבוד למשימות שעשויות לחרוג מפרק זמן לעיבוד של הודעת ברירת המחדל למשך 10 שניות. אם לא מאריכים את תקופת ההחזקה של ההודעות, יכול להיות שהן יישלחו מחדש ושהן יעברו עיבוד כפול.
  • טיפול בשגיאות: תורים של Spanner מנסים לשלוח מחדש הודעות שעיבוד שלהן נכשל עם השהיה אקספוננציאלית במהלך השעה הראשונה, והודעות ישנות יותר יישלחו מחדש פעם בשעה. כדאי להעביר הודעות שנכשלו באופן קבוע לתור נפרד.
  • מעקב: עוקבים אחרי עומק התור והגיל של ההודעה הכי ישנה שלא אושרה כדי לוודא שהמקבלים לא מגבילים את קצב ההעברה של הנתונים.
  • אידמפוטנטיות: כדאי לתכנן את מעבדי ההודעות כך שיאפשרו חזרה על פעולות, כי מסירה אחת לפחות פירושה שאפשר לקבל מדי פעם כפילויות. מידע נוסף זמין בדף בנושא עיבוד של כל הודעה פעם אחת בלבד ואישור של כל הודעה פעם אחת לכל היותר.
  • משך הזמן של פונקציה שמחזירה טבלה: מומלץ להימנע ממשכי זמן קצרים מדי או ארוכים מדי של פונקציות שמחזירות טבלה. מומלץ להגדיר משך זמן מתון, כמו 20 דקות.
  • גודל אצווה: כדאי להתאים את max_batch_size לעומס העבודה. כדי להימנע ממחלוקת על נעילה בשורות משותפות, כדאי להשתמש באצוות קטנות יותר לאירועים עם fan-out גבוה. כדאי להשתמש באצוות גדולות יותר למשימות עצמאיות עם מספר גבוה של שאילתות לשנייה. כדי להשיג את הביצועים הטובים ביותר, מומלץ לאשר או להאריך את תקופת ההחזקה של הודעות בקבוצה בעסקה אחת.

מעקב

אפשר לעקוב אחרי פעולות בתור באמצעות טבלאות של Spanner introspection. הטבלאות האלה לא כוללות עמודות ספציפיות לתור, אבל אפשר לזהות פעילות בתור על ידי חיפוש של שמות התורים שהוגדרו על ידי המשתמש בטבלאות הבאות:

מדדי התור הבאים של Spanner נמצאים תחת הקידומת spanner.googleapis.com/queue/* ב-Cloud Monitoring:

  • buffered_ready_messages: (GAUGE, INT64, 1) מספר ההודעות שמוחזקות בזיכרון ומוכנות למסירה לנמען.
  • message_send_count: (DELTA, INT64, 1) מספר ההודעות שנשלחו ב-Spanner במהלך המרווח לתור.
  • message_ack_count: (DELTA, INT64, 1) מספר ההודעות שאושרה קבלתן ב-Spanner במהלך המרווח עבור תור.
  • oldest_unacked_message_age: (GAUGE, INT64, 1) הגיל (בשניות) של ההודעה הכי ישנה בתור שלא נשלח לגביה אישור.
  • lease_expiration_count: (DELTA, INT64, 1) מספר התפוגות של ההסכם על השכרה ב-Spanner במהלך המרווח של תור.

כל המדדים הקודמים נדגמים בערך כל 60 שניות. אחרי הדגימה, יכול להיות שיעברו עד 120 שניות עד שהנתונים יוצגו. רישום ביומן הביקורת של Spanner כולל כתיבות, קריאות ופעולות סכימה בתורים.

פתרון בעיות

בקטעים הבאים מוסבר איך לזהות ולפתור בעיות נפוצות בשימוש בתורים של Spanner.

הצטברות של הודעות שלא נשלחו

אבחון

הגישה ל-oldest_unacked_message_age ול-buffered_ready_messages היא גישה ברמת אדמין. המשמעות היא שיש חוסר איזון בין קצב שליחת ההודעות לבין יכולת עיבוד ההודעות של האפליקציה.

רזולוציה

כדי לפתור את הבעיה:

  • בודקים את שיעור האישור: אם המדד message_ack_count ירד, בודקים את עובדי הלקוח כדי לוודא שהם פועלים בצורה תקינה ולא נתקעו או קרסו.
  • בודקים את קצב השליחה: אם יש עלייה חדה בערך של message_send_count, צריך להגדיל את מספר העובדים לעיבוד הודעות כדי להתמודד עם העומס המוגבר.
  • מוודאים שנגמרו המשאבים: בודקים אם הערכים של buffered_ready_messages ושל lease_expiration_count גבוהים. השילוב הזה מצביע על מחסור בפונקציות פעילות שמחזירות ערך של טבלה (TVF) או על עיבוד איטי של הלקוח.

הודעות פרטיות נתקעות

אבחון

המדד oldest_unacked_message_age גבוה, אבל המדד buffered_ready_messages נמוך או יציב. המשמעות היא שעיבוד ההודעות הספציפיות נכשל או שהן לא אושרו, ולא שיש צוואר בקבוק בקיבולת הכוללת.

רזולוציה

כדי לפתור את הבעיה:

  • זיהוי הודעות תקועות: שולחים שאילתה לטבלת התור ומסדרים את התוצאות לפי זמן המסירה כדי למצוא את ההודעות הכי ישנות שלא אושרו:

    GoogleSQL

    SELECT *
    FROM UserTasks
    ORDER BY DeliverTime ASC
    LIMIT 10;
    

    PostgreSQL

    SELECT *
    FROM usertasks
    ORDER BY deliver_time ASC
    LIMIT 10;
    
  • בדיקת כשלים בעיבוד: בודקים את יומני האפליקציה כדי להבין למה העובדים לא מאשרים את ההודעות שזוהו. הודעות שלא אושרו נשלחות מחדש באופן אוטומטי אחרי שתוקף ההרשאה שלהן פג.

תוקף ההודעות פג לפני שהעיבוד מסתיים

אבחון

המדד lease_expiration_count גבוה או בעלייה. המשמעות היא שזמן העיבוד של ההודעה חורג ממשך ההקצאה (ברירת מחדל היא 10 שניות) לפני שהעובדים יכולים לאשר את ההודעות.

רזולוציה

כדי לפתור את הבעיה:

  • חידוש חוזי שכירות באופן יזום: אם עיבוד ההודעה נמשך יותר מ-10 שניות, צריך להתקשר אל RENEWLEASE_QUEUE_NAME() TVF באופן תקופתי. כדי להתחשב בבטחה בזמן האחזור של הרשת, צריך לחדש את ההקצאה בערך 7-8 שניות אחרי תחילת העיבוד.
  • בדיקה של עיבוד איטי: אם האפליקציה כבר מחדשת את ההשכרות באופן פעיל, אבל lease_expiration_count עדיין גבוה, כדאי לבדוק את קוד ה-Backend כדי לזהות צווארי בקבוק בעיבוד, קריאות RPC איטיות או חסימות הדדיות.

אי אפשר לשנות את קצב שליחת ההודעות

אבחון

המדד message_send_count מגיע לרמת תפוקה מקסימלית או שבקשות הפרסום נתקלות בחביון גבוה של כתיבה כשמנסים להגדיל את קצב השליחה.

רזולוציה

כדי לפתור את הבעיה, כדאי לבצע את השינויים המבניים הבאים:

  • הגדלת משאבי המחשוב: מוסיפים צמתים או יחידות עיבוד למופע Spanner כדי להגדיל את הקיבולת הכוללת של מסד הנתונים.
  • בדיקת פיצולים: כדאי להעריך אם אפשר להוסיף עוד פיצולים כדי לחלק את עומס הכתיבה בין כמה שרתים.
  • אופטימיזציה של הניתוב: מוודאים שאפליקציית הפרסום כותבת ישירות לאזור הראשי של מופע Spanner כדי לצמצם את זמן האחזור של הכתיבה.

הודעות מוכנות לא מעובדות

אבחון

הערך של המדד buffered_ready_messages גבוה ועולה. המשמעות היא שההודעות נשמרות בזיכרון ומוכנות למסירה, אבל תהליכי העבודה של המקבל לא שולפים אותן.

רזולוציה

כדי לפתור את הבעיה:

  • בדיקת חיבורי TVF פעילים: בודקים את מספר חיבורי TVF הפעילים של המקלט כדי לוודא שעובדי הקריאה מתחברים באופן פעיל ושולפים הודעות. אם העובדים התנתקו או שלא מופעלות מספיק שאילתות בו-זמניותRECEIVE_QUEUE_NAME(), ההודעות נשארות במאגר הזמני ללא בדיקה.

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