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

자바

dbClient.write(
        Collections.singletonList(
            Mutation.newSendBuilder("UserTasks")
                .setKey(Key.of(123L, "some-unique-id-1"))
                .setPayload(Value.bytes(ByteArray.copyFrom("message3")))
                .setDeliveryTime(futureTime)
                .build()));

메시지 수신하기

ExecuteStreamingSQL와 함께 RECEIVE_QUEUE_NAME() 테이블 값 함수 (TVF)를 사용하여 메시지를 수신합니다. 이는 장기 실행 호출입니다. 작업자당, 대기열당 이러한 호출 중 하나를 루프 방식으로 실행해야 합니다. 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}),
})

자바

dbClient.write(
    Collections.singletonList(
        Mutation.newAckBuilder("UserTasks")
            .setKey(Key.of(2L))
            .build()));

권장사항

다음은 Spanner 대기열 사용에 관한 권장사항입니다.

  • 작은 페이로드: 대기열 메시지 페이로드를 4KB 미만으로 작게 유지합니다. 더 큰 데이터에는 대역 외 저장소 패턴을 사용합니다.
  • 임대 관리: 기본 메시지 임대 기간인 10초를 초과할 수 있는 작업의 임대 기간을 연장합니다. 메시지 리스를 연장하지 않으면 재전송이 발생하고 이중 처리가 발생할 수 있습니다.
  • 오류 처리: Spanner는 첫 시간 내에 처리에 실패한 재시도 메시지를 지수 백오프를 사용하여 대기열에 추가하며, 오래된 메시지는 시간당 한 번 재시도됩니다. 영구적으로 실패한 메시지를 별도의 큐로 이동하는 것이 좋습니다.
  • 모니터링: 수신자가 파이프라인 인입을 제한하지 않도록 대기열 깊이와 확인되지 않은 가장 오래된 메시지 기간을 모니터링합니다.
  • 멱등성: 최소 1회 전송은 간혹 중복이 발생할 수 있음을 의미하므로 반복 작업을 허용하도록 메시지 프로세서를 설계하세요. 자세한 내용은 정확히 한 번 처리 및 최대 한 번 확인 페이지를 참고하세요.
  • 테이블 값 함수 기간: 테이블 값 함수의 기간이 너무 짧거나 너무 길지 않도록 합니다. 20분과 같은 적당한 시간이 권장됩니다.
  • 배치 크기 조정: 워크로드에 맞게 max_batch_size를 조정합니다. 공유 행에서 잠금 경합을 방지하려면 팬아웃이 많은 이벤트에 더 작은 일괄 처리를 사용하세요. 독립적이고 초당 쿼리 수가 많은 작업에는 더 큰 배치 크기를 사용합니다. 최적의 성능을 위해 단일 트랜잭션에서 배치에 대한 메시지 임대를 확인하거나 연장합니다.

모니터링

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() 쿼리를 실행하지 않으면 메시지가 버퍼에서 폴링되지 않은 상태로 유지됩니다.

다음 단계