YQL-запросы к топикам

Для чтения и записи сообщений в топики используются привычные YQL-конструкции: SELECT для чтения и INSERT для записи.

Локальные и внешние топики

YQL-запросы к топикам работают одинаково независимо от того, находится топик в текущей базе или в другой базе YDB. Источником и приёмником сообщений может быть как топик в той же базе данных, в которой выполняется запрос, так и топик в другой базе.

Локальные топики

Локальные топики — топики, созданные в той же базе YDB, что и выполняемый запрос.

В тексте запроса к ним обращаются по короткому имени — так же, как к таблице в текущей базе:

SELECT * FROM input_topic WITH (FORMAT = json_each_row, SCHEMA = (...));
INSERT INTO output_topic SELECT ...;

Внешние топики

Внешние топики — топики, расположенные в другой базе YDB.

Доступ к ним выполняется только через заранее созданный внешний источник данных с типом источника YDB.

После создания источника, например с именем ext_source, обращение к топику input_topic во внешней базе записывается так:

SELECT * FROM ext_source.input_topic WITH (FORMAT = json_each_row, SCHEMA = (...));

Имя ext_source в документации условное — в вашей базе источник может называться иначе; важно, чтобы оно совпадало в CREATE EXTERNAL DATA SOURCE и в префиксе перед именем топика.

Чтение из топика

Чтение из топика можно выполнять в табличном и потоковом режимах (не путать с потоковыми запросами).

Табличное чтение

В табличном режиме чтение выполняется от первого до последнего смещения, хранящегося в топике на момент запуска запроса. Если в топик продолжается запись данных, то запрос остановится после достижения последнего смещения, известного на момент запуска. Указание фильтров по Служебным полям ускоряет чтение, так как чтение происходит только по указанным диапазонам.

SELECT
    Data    -- тело сообщения
FROM
    input_topic  -- локальный топик; для внешнего: ext_source.input_topic
LIMIT 10;

Потоковое чтение

Для чтения новых сообщений используйте опцию WITH (STREAMING = "TRUE") — подробнее в разделе Потоковое чтение данных из топика. Чтение начинается с текущего момента и продолжается до тех пор, пока не будет прочитано заданное в выражении LIMIT количество сообщений. Параметр LIMIT обязателен — без него запрос не завершится, так как будет ожидать новые сообщения бесконечно.

SELECT
    Data
FROM
    ext_source.input_topic  -- внешний топик; для локального: input_topic
WITH (STREAMING = "TRUE")
LIMIT 10;

Для непрерывной обработки поступающих данных используйте потоковые запросы.

Формат и схема сообщений

При чтении из топика тело сообщения можно получить двумя способами: сырые данные и форматированные данные.

Сырые данные

Используйте, когда содержимое сообщения не нужно разбирать — достаточно прочитать тело как есть.

SELECT
    Data
FROM
    input_topic  -- локальный топик; для внешнего: ext_source.input_topic
WITH (
    FORMAT = raw,
    SCHEMA = (
        Data String
    )
)
LIMIT 10;

В результате доступна только колонка Data — тело сообщения в исходном виде.

Тот же результат можно получить без блока WITH — см. табличное чтение.

Форматированные данные

Используйте, когда сообщения сериализованы в известном формате (JSON, CSV и др.). Параметр FORMAT задаёт способ разбора, а SCHEMA — имена и типы полей, которые появятся в результате SELECT:

SELECT
    Id,
    Name
FROM
    input_topic  -- локальный топик; для внешнего: ext_source.input_topic
WITH (
    FORMAT = json_each_row,
    SCHEMA = (
        Id Uint64 NOT NULL,
        Name Utf8 NOT NULL
    )
);

Поля из SCHEMA доступны в SELECT по имени — как колонки таблицы.

Подробнее о поддерживаемых форматах: Форматы данных при чтении/записи из топиков.

Использование читателя

Читатель (consumer) — именованная подписка на топик, которая хранит текущую позицию чтения.

Читатель создаётся через CLI или при создании топика с помощью CREATE TOPIC. Имя читателя указывается в тексте запроса прагмой:

PRAGMA pq.Consumer="my_consumer";

Если читатель не указан, чтение из топика выполняется без него. Указание читателя позволяет отслеживать позицию чтения и лаг со стороны топика, например через CLI.

Перенос данных из топика в таблицу через UPSERT

Данные из топика можно переложить в таблицу через UPSERT INTO:

UPSERT INTO
    table_name
SELECT
    Data            -- можно использовать любые преобразования
FROM
    ext_source.input_topic;  -- внешний топик; для локального: input_topic

Служебные поля

При чтении можно запрашивать служебные поля:

SELECT
    Data,                                                   -- тело сообщения
    __ydb_create_time AS CreateTime,                        -- время создания сообщения
    __ydb_write_time AS WriteTime,                          -- время записи сообщения
    __ydb_offset AS Offset,                                 -- смещение сообщения в топике
    __ydb_partition_id AS Partition,                        -- номер партиции
    __ydb_message_group_id AS MessageGroupId,               -- идентификатор группы сообщений
    __ydb_seq_no AS SeqNo                                   -- порядковый номер внутри партиции
FROM
    input_topic  -- локальный топик; для внешнего: ext_source.input_topic
LIMIT 10;

Фильтры по служебным полям вычисляются до чтения данных из топика и существенно сокращают объём считываемых сообщений. Поддерживаются операторы сравнения (=, <>, <, <=, >, >=, IN), логические условия (AND, OR) и поля partition_id, write_time, offset. Предикаты по остальным служебным полям не ограничивают объём чтения.

SELECT
    Data
FROM
    ext_source.input_topic  -- внешний топик; для локального: input_topic
WHERE
    __ydb_partition_id = 42
        AND __ydb_offset >= 1000
        AND __ydb_offset <= 1100
        AND __ydb_write_time > CurrentUtcTimestamp() - Interval("PT2H");

Запись в топик

Запись одного сообщения

INSERT INTO
    output_topic  -- локальный топик; для внешнего: ext_source.output_topic
SELECT
    "my_message";       -- тело сообщения

Запись данных из таблицы

Чтобы записать в топик строку таблицы с несколькими колонками, сформируйте JSON-объект. Функция TableRow создаёт структуру из всех колонок, Yson::From преобразует её в Yson, Yson::SerializeJson сериализует в JSON-строку, а ToBytes конвертирует результат в тип String, который требуется для записи в топик:

INSERT INTO
    ext_source.output_topic  -- внешний топик; для локального: output_topic
SELECT
    ToBytes(Unwrap(Yson::SerializeJson(Yson::From(TableRow()))))
FROM
    table_name;

Ограничения

Важно

Чтение и запись пользовательских атрибутов не поддерживаются.

Важно

Транзакционная запись через YQL/INSERT INTO не поддерживается — в топике могут появиться частичные результаты запроса.

См. также

Предыдущая
Следующая