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 ステートメントで明示的に作成する必要がない唯一の列は、GoogleSQL では DeliverTime、PostgreSQL では deliver_time と呼ばれます。これらは 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 ミューテーションを使用してメッセージを挿入します。

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

メールの受信

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 ミューテーションを使用して、メッセージの確認応答を行います。

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

ベスト プラクティス

Spanner キューを使用する場合のベスト プラクティスは次のとおりです。

  • 小さいペイロード: キュー メッセージ ペイロードを 4 KB 未満に抑えます。大きなデータには、アウトオブバンド ストレージ パターンを使用します。
  • リースの管理: デフォルトのメッセージ リース(10 秒)を超える可能性のあるタスクのリースを延長します。メッセージ リースを延長できないと、再配信や二重処理が発生する可能性があります。
  • エラー処理: Spanner キューは、処理に失敗したメッセージを最初の 1 時間以内にバックオフで再試行します。古いメッセージは 1 時間に 1 回再試行されます。永続的に失敗したメッセージを別のキューに移動することを検討してください。
  • モニタリング: キューの深さと最も古い未確認メッセージの経過時間をモニタリングして、受信側がパイプラインの取り込みを制限していないことを確認します。
  • べき等性: at-least-once 配信では重複が発生する可能性があるため、繰り返しオペレーションを許容するようにメッセージ プロセッサを設計します。詳細については、1 回限りの処理と 1 回限りの確認応答のページをご覧ください。
  • テーブル値関数の期間: テーブル値関数の期間が極端に短い場合や、極端に長い場合は避けてください。20 分程度の適度な時間をおすすめします。
  • バッチサイズ設定: ワークロードに合わせて max_batch_size を調整します。ファンアウトの多いイベントには、共有行でのロック競合を避けるために、バッチサイズを小さくします。独立した 1 秒あたりのクエリ数の多いタスクには、大きなバッチを使用します。最高のパフォーマンスを得るには、単一のトランザクションでバッチのメッセージ リースを確認または延長します。

モニタリング

Spanner のイントロスペクション テーブルを使用して、キュー オペレーションをモニタリングできます。これらのテーブルにはキュー固有の列は含まれていませんが、次のテーブルでユーザー定義のキュー名を検索することで、キュー アクティビティを特定できます。

次の Spanner キュー指標は、Cloud Monitoringspanner.googleapis.com/queue/* 接頭辞にあります。

  • 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_agebuffered_ready_messages の両方が昇格している。これは、メッセージ送信レートとアプリケーションのメッセージ処理能力の間に不均衡があることを示します。

解決策

この問題を解決するには、次の操作を行います。

  • 確認応答率を確認する: message_ack_count 指標が低下した場合は、クライアント ワーカーが正しく実行され、停止またはクラッシュしていないことを確認します。
  • 送信レートを確認します。message_send_count が急増した場合は、メッセージ処理ワーカーをスケールアップして、増加した負荷を処理します。
  • リソースの枯渇を確認する: buffered_ready_messageslease_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 が高い場合は、バックエンド コードで処理のボトルネック、RPC 呼び出しの遅延、デッドロックがないか確認します。

メッセージ送信レートをスケーリングできない

診断

送信レートを上げようとすると、message_send_count 指標がスループットのプラトーに達するか、パブリッシュ リクエストで書き込みレイテンシが増加します。

解決策

この問題を解決するには、次の構造変更を検討してください。

  • コンピューティング リソースをスケールアップする: Spanner インスタンスにノードまたは処理ユニットを追加して、データベースの全体的な容量を増やします。
  • スプリットを確認する: 書き込み負荷を複数のサーバーに分散するために、スプリットを追加できるかどうかを評価します。
  • ルーティングを最適化する: パブリッシュ アプリケーションが Spanner インスタンスのリーダー リージョンに直接書き込み、書き込みレイテンシを最小限に抑えるようにします。

準備完了メッセージが処理されない

診断

buffered_ready_messages 指標が高く、増加している。これは、メッセージがバッファリングされ、メモリ内の配信準備が整っているが、受信側ワーカーがメッセージを取得していないことを示します。

解決策

この問題を解決するには、次の操作を行います。

  • アクティブな TVF 接続を確認する: アクティブなレシーバー TVF の数を確認して、リーダー ワーカーがアクティブに接続してメッセージをプルしていることを確認します。ワーカーが切断された場合や、十分な数の同時 RECEIVE_QUEUE_NAME() クエリが実行されていない場合、メッセージはバッファ内でポーリングされないままになります。

次のステップ