Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Движок таблицы AzureQueue

Этот движок интегрируется с экосистемой 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. Для этого:

  1. С помощью движка создайте таблицу для чтения из указанного пути в Azure Blob Storage и рассматривайте её как поток данных.
  2. Создайте таблицу с нужной структурой.
  3. Создайте 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, однако есть несколько важных отличий:

  1. Используйте system.azure_queue_metadata_cache для состояния очереди в памяти в версиях сервера >= 25.1. Для более старых версий используйте system.s3queue_metadata_cache (он также содержит информацию о таблицах azure).
  2. Используйте таблицу system.azure_queue_metadata, чтобы напрямую проверить состояние, хранящееся в Keeper: количество узлов processed, processing и failed для каждого объекта метаданных, а при необходимости — их содержимое. Это аналог system.s3_queue_metadata для AzureQueue.
  3. Включите 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 = 1fsync_part_directory = 1) для целевой таблицы MergeTree существенно сокращает это окно.

Navigation