Этот движок интегрируется с экосистемой Azure Blob Storage, обеспечивая потоковый импорт данных.
CREATE TABLE
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 Blob Storage для перемещения успешно обработанных файлов, если пункт назначения — другой контейнер Azure.
Возможные значения:
- 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';SELECT из таблиц на движке таблицы AzureQueue
Запросы SELECT для таблиц AzureQueue по умолчанию запрещены. Это соответствует распространённому паттерну очереди, при котором данные считываются один раз, а затем удаляются из очереди. SELECT запрещён, чтобы предотвратить случайную потерю данных.
Однако в некоторых случаях он может быть полезен. Для этого нужно установить значение настройки stream_like_engine_allow_direct_select в True.
У движка AzureQueue есть специальная настройка для запросов SELECT: commit_on_select. Установите для неё значение False, чтобы сохранить данные в очереди после чтения, или True, чтобы удалить их.
Описание
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;Виртуальные столбцы
_path— путь к файлу._file— имя файла.
Подробнее о виртуальных столбцах см. здесь.
Интроспекция
Включите логирование для таблицы с помощью настройки таблицы enable_logging_to_queue_log=1.
Возможности интроспекции такие же, как у движка таблицы S3Queue, однако есть несколько важных отличий:
- Используйте
system.azure_queue_metadata_cacheдля состояния очереди в памяти в версиях сервера >= 25.1. Для более старых версий используйтеsystem.s3queue_metadata_cache(он также содержит информацию о таблицахazure). - Используйте таблицу
system.azure_queue_metadata, чтобы напрямую проверить состояние, хранящееся в Keeper: количество узловprocessed,processingиfailedдля каждого объекта метаданных, а при необходимости — их содержимое. Это аналогsystem.s3_queue_metadataдляAzureQueue. - Включите
system.azure_queue_logчерез основную конфигурацию ClickHouse, например:
<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, может незаметно привести к потере обработанных строк: файл помечается в Keeper как обработанный (а при after_processing = 'delete' также удаляется исходный blob) сразу после завершения вставки, однако вставленные строки надёжно сохраняются только после выполнения fsync целевой части, который по умолчанию не выполняется синхронно (fsync_after_insert = 0). Для рекомендуемого способа обработки через материализованное представление установка fsync_after_insert = 1 (и fsync_part_directory = 1) для целевой таблицы MergeTree существенно сокращает это окно.