AzureQueue 테이블 엔진

AzureQueue 테이블 엔진

Azure Blob Storage 생태계와의 통합을 제공해, 스트리밍 데이터 가져오기를 허용하는 테이블 엔진이에요.

출처: 문서

본문

이 엔진은 Azure Blob Storage 생태계와의 통합을 제공해, 스트리밍 데이터 가져오기를 허용해요.

테이블 만들기 (Create table)

CREATE TABLE test (name String, value UInt32)
    ENGINE = AzureQueue(...)
    [SETTINGS]
    [mode = '',]
    [after_processing = 'keep',]
    [keeper_path = '',]
    ...

엔진 매개변수 (Engine parameters)

AzureQueue 매개변수는 AzureBlobStorage 테이블 엔진이 지원하는 것과 같아요. 매개변수 섹션은 여기를 참고해요. AzureBlobStorage 테이블 엔진과 유사하게, 사용자는 로컬 Azure Storage 개발에 Azurite 에뮬레이터를 사용할 수 있어요. 자세한 내용은 여기를 참고해요.

예시 (Example)

CREATE TABLE azure_queue_engine_table
(
    `key` UInt64,
    `data` String
)
ENGINE = AzureQueue('DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://azurite1:10000/devstoreaccount1/;', 'testcontainer', '*', 'CSV')
SETTINGS mode = 'unordered'

설정 (Settings)

지원되는 설정 집합은 대부분 S3Queue 테이블 엔진과 같지만 s3queue_ 접두사가 없어요. 전체 설정 목록을 참고해요. 테이블에 대해 구성된 설정 목록을 얻으려면 system.azure_queue_settings 테이블을 사용해요. 24.10부터 사용 가능해요. 아래는 AzureQueue에만 호환되고 S3Queue에는 적용되지 않는 설정들이에요.

after_processing_move_connection_string

성공적으로 처리된 파일을 이동할(대상이 다른 Azure 컨테이너인 경우) Azure Blob Storage의 연결 문자열. 가능한 값: 문자열. 기본값: 빈 문자열.

after_processing_move_container

성공적으로 처리된 파일을 이동할(대상이 다른 Azure 컨테이너인 경우) 컨테이너 이름. 가능한 값: 문자열. 기본값: 빈 문자열.

예시:

CREATE TABLE azure_queue_engine_table
(
    `key` UInt64,
    `data` String
)
ENGINE = AzureQueue('DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://azurite1:10000/devstoreaccount1/;', 'testcontainer', '*', 'CSV')
SETTINGS
    mode = 'unordered',
    after_processing = 'move',
    after_processing_move_connection_string = 'DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://azurite1:10000/devstoreaccount1/;',
    after_processing_move_container = 'dst-container';

AzureQueue 테이블 엔진에서 SELECT

AzureQueue 테이블에서 SELECT 쿼리는 기본적으로 금지돼요. 이것은 데이터가 한 번 읽힌 후 큐에서 제거되는 일반적인 큐 패턴을 따르는 것이에요. 우발적 데이터 손실을 막기 위해 SELECT가 금지돼 있어요. 하지만 때로 유용할 수 있어요. 이렇게 하려면 설정 stream_like_engine_allow_direct_selectTrue로 설정해야 해요. AzureQueue 엔진에는 SELECT 쿼리를 위한 특별한 설정이 있어요. commit_on_select. 읽은 후 큐에 데이터를 보존하려면 False, 제거하려면 True로 설정해요. (참고: 이 설정은 exclusive 모드에서는 의미가 없고 무시돼요. exclusive 모드는 항상 commit_on_selectTrue인 것처럼 동작해요.)

설명 (Description)

SELECT는 각 파일을 한 번만 가져올 수 있으므로(디버깅 외에는) 스트리밍 가져오기에 특히 유용하지 않아요. materialized views로 실시간 스레드를 만드는 것이 더 실용적이에요. 이렇게 하려면:

  • 엔진을 사용해 Azure Blob Storage의 지정된 경로에서 소비하는 테이블을 만들고 그것을 데이터 스트림으로 고려해요.
  • 원하는 구조의 테이블을 만들어요.
  • 엔진에서 데이터를 변환해 이전에 만든 테이블에 넣는 materialized view를 만들어요.

MATERIALIZED VIEW가 엔진에 결합되면 백그라운드에서 데이터 수집을 시작해요. 엔진 인자의 형태는 AzureQueue(connection_string, container_name, blobpath, format[, compression])이에요. 예시:

CREATE TABLE azure_queue_engine_table (key UInt64, data String)
ENGINE=AzureQueue('DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://azurite1:10000/devstoreaccount1/;', 'testcontainer', '*', 'CSV')
SETTINGS
      mode = 'unordered';

CREATE TABLE stats (key UInt64, data String)
ENGINE = MergeTree() ORDER BY key;

CREATE MATERIALIZED VIEW consumer TO stats
AS SELECT key, data FROM azure_queue_engine_table;

SELECT * FROM stats ORDER BY key;

가상 컬럼 (Virtual columns)

  • _path — 파일 경로.
  • _file — 파일 이름.

가상 컬럼에 대한 자세한 정보는 여기를 참고해요.

내부 검사 (Introspection)

테이블 설정 enable_logging_to_queue_log=1로 테이블에 로깅을 활성화해요. 내부 검사 기능은 S3Queue 테이블 엔진과 같고 몇 가지 뚜렷한 차이가 있어요.

  • 서버 버전 >= 25.1에서는 큐의 인메모리 상태에 system.azure_queue_metadata_cache를 사용해요. 이전 버전에서는 system.s3queue_metadata_cache를 사용해요(azure 테이블의 정보도 담아요).
  • system.azure_queue_metadata 테이블로 keeper에 저장된 상태를 직접 검사해요. 메타데이터 객체당 processed, processing, failed 노드 수와, 필요 시 그 내용을 보여줘요. 이것은 system.s3_queue_metadata의 AzureQueue 대응물이에요.
  • 메인 ClickHouse 설정으로 system.azure_queue_log를 활성화해요. 예:
<azure_queue_log>
    <database>system</database>
    <table>azure_queue_log</table>
</azure_queue_log>

이 영구 테이블은 system.s3queue_metadata_cache와 같은 정보를 담지만 처리되고 실패한 파일에 대한 것이에요. 테이블의 구조는 다음과 같아요.

CREATE TABLE system.azure_queue_log
(
    `hostname` LowCardinality(String) COMMENT 'Hostname',
    `event_date` Date COMMENT 'Event date of writing this log row',
    `event_time` DateTime COMMENT 'Event time of writing this log row',
    `database` String COMMENT 'The name of a database where current S3Queue table lives.',
    `table` String COMMENT 'The name of S3Queue table.',
    `uuid` String COMMENT 'The UUID of S3Queue table',
    `file_name` String COMMENT 'File name of the processing file',
    `rows_processed` UInt64 COMMENT 'Number of processed rows',
    `status` Enum8('Processed' = 0, 'Failed' = 1) COMMENT 'Status of the processing file',
    `processing_start_time` Nullable(DateTime) COMMENT 'Time of the start of processing the file',
    `processing_end_time` Nullable(DateTime) COMMENT 'Time of the end of processing the file',
    `exception` String COMMENT 'Exception message if happened'
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_time)
COMMENT 'Contains logging entries with the information files processes by S3Queue engine.'

예시:

SELECT *
FROM system.azure_queue_log
LIMIT 1
FORMAT Vertical

Row 1:
──────
hostname:              clickhouse
event_date:            2024-12-16
event_time:            2024-12-16 13:42:47
database:              default
table:                 azure_queue_engine_table
uuid:                  1bc52858-00c0-420d-8d03-ac3f189f27c8
file_name:             test_1.csv
rows_processed:        3
status:                Processed
processing_start_time: 2024-12-16 13:42:47
processing_end_time:   2024-12-16 13:42:47
exception:

1 row in set. Elapsed: 0.002 sec.

제한 (Limitations)

AzureQueueS3Queue와 같은 구현을 공유하며 같은 제한을 가져요. 특히 ClickHouse 노드의 장치 수준 전원 손실은 소비된 행을 조용히 잃을 수 있어요. insert가 끝나는 즉시 파일이 Keeper에서 처리된 것으로 기록되고(after_processing = 'delete'면 소스 blob이 제거됨) 삽입된 행은 대상 파트가 fsync된 후에야 영속적이 되는데, 이는 기본적으로 동기적으로 일어나지 않기 때문이에요(fsync_after_insert = 0). 권장되는 materialized-view 소비 경로에서는 대상 MergeTree 테이블에 fsync_after_insert = 1(그리고 fsync_part_directory = 1)을 설정하면 이 창이 상당히 줄어요.

더 알아보기 (Learn more)