このエンジンは Azure Blob Storage エコシステムと連携し、ストリーミングデータをインポートできます。
テーブルの作成
CREATE TABLE test (name String, value UInt32)
ENGINE = AzureQueue(...)
[SETTINGS]
[mode = '',]
[after_processing = 'keep',]
[keeper_path = '',]
...エンジンパラメータ
AzureQueue のパラメータは、AzureBlobStorage テーブルエンジンでサポートされているものと同じです。パラメータについては、こちらを参照してください。
AzureBlobStorage テーブルエンジンと同様に、ローカルで Azure Storage を開発する際には Azurite エミュレータを使用できます。詳細はこちらを参照してください。
例
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'設定
サポートされている設定のほとんどは S3Queue テーブルエンジンと同じですが、s3queue_ プレフィックスは付きません。設定の完全な一覧を参照してください。
テーブルに対して設定されている項目の一覧を取得するには、system.azure_queue_settings テーブルを使用します。24.10 から利用できます。
以下は、AzureQueue でのみ使用でき、S3Queue には適用されない設定です。
after_processing_move_connection_string
宛先が別の Azure コンテナーである場合に、正常に処理されたファイルの移動先として使用する Azure Blob Storage の接続文字列。
設定可能な値:
- String。
デフォルト値: 空文字列。
after_processing_move_container
宛先が別の Azure コンテナーである場合に、正常に処理されたファイルの移動先となるコンテナーの名前。
設定可能な値:
- String.
デフォルト値: 空文字列。
例:
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 が禁止されているのは、意図しないデータ損失を防ぐためです。
ただし、場合によっては直接 SELECT したいこともあります。その場合は、設定 stream_like_engine_allow_direct_select を True にする必要があります。
AzureQueue engine には、SELECT クエリ向けの特別な設定 commit_on_select があります。読み取り後もキュー内のデータを保持するには False に設定し、削除するには True に設定します。
説明
SELECT はストリーミングインポートにはあまり適していません (デバッグ用途を除く) 。各ファイルは一度しかインポートできないためです。代わりに、materialized view を使ってリアルタイムの処理パイプラインを作成するほうが実用的です。手順は次のとおりです。
- エンジンを使用して、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;仮想カラム
_path— ファイルのパス。_file— ファイル名。
仮想カラムの詳細については、こちらを参照してください。
イントロスペクション
テーブル設定 enable_logging_to_queue_log=1 で、このテーブルのログ記録を有効にします。
イントロスペクション機能は S3Queue テーブルエンジン と同じですが、いくつか異なる点があります。
- server バージョンが >= 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:
─制限事項
AzureQueue は S3Queue と同じ実装を使用しており、同じ制限事項があります。特に、ClickHouse ノードでデバイスレベルの電源障害が発生すると、消費済みの行が気付かないうちに失われる可能性があります。insert が完了するとすぐに、ファイルは Keeper に処理済みとして記録され (after_processing = 'delete' の場合はソースブロブも削除されます) 、一方で挿入された行が永続化されるのはターゲット part が fsync された後です。これはデフォルトでは同期的に行われません (fsync_after_insert = 0) 。推奨されるマテリアライズドビューを介した消費パスでは、ターゲットの MergeTree テーブルで fsync_after_insert = 1 (および fsync_part_directory = 1) を設定することで、この時間的なギャップを大幅に縮小できます。