Типичные шаблоны потоковых запросов

В этом разделе собраны минимальные примеры потоковых запросов для наиболее распространённых сценариев. Сначала описан базовый шаблон чтения данных из топика, затем — варианты полноценной работы с данными: обработка данных и запись результатов в топик в формате JSON, в топик в виде строки и в таблицу. Каждый пример можно использовать как отправную точку для собственных задач.

В примерах ниже используются локальные и внешние топики. Обозначения:

  • ext_source — заранее созданный external data source;
  • input_topic — топик, откуда читаются данные;
  • output_topic — топик, куда записываются результаты;
  • output_table — таблица YDB, куда записываются результаты.

Чтение данных из топика

Чтение структурированных сообщений выполняется с помощью SELECT ... FROM ... WITH (FORMAT, SCHEMA). Блок WITH задаёт формат входных данных и схему — какие поля ожидаются в каждом сообщении и их типы.

Следующий фрагмент используется внутри CREATE STREAMING QUERY в блоке DO BEGIN ... END DO:

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

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

Этот шаблон используется во всех последующих примерах.

Запись в топик (JSON)

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

CREATE STREAMING QUERY write_json_example AS
DO BEGIN

INSERT INTO ext_source.output_topic -- или локальный топик output_topic
SELECT
    -- Формирование JSON из отдельных полей
    ToBytes(Unwrap(Yson::SerializeJson(Yson::From(
        AsStruct(Id AS id, Name AS name)
    ))))
FROM
    input_topic -- или внешний топик ext_source.input_topic
WITH (
    FORMAT = json_each_row,  -- Формат входных данных
    SCHEMA = (               -- Схема входных данных
        Id Uint64 NOT NULL,
        Name Utf8 NOT NULL
    )
);

END DO

Подробнее о функциях:

Запись в топик (строка)

Запрос читает события из входного топика и записывает в выходной топик одно поле в виде строки. Для записи в топик строк SELECT должен возвращать одну колонку типа String или Utf8.

CREATE STREAMING QUERY write_utf8_example AS
DO BEGIN

INSERT INTO output_topic -- или внешний топик ext_source.output_topic
SELECT
    Name
FROM
    ext_source.input_topic -- или локальный топик input_topic
WITH (
    FORMAT = json_each_row,  -- Формат входных данных
    SCHEMA = (               -- Схема входных данных
        Id Uint64 NOT NULL,
        Name Utf8 NOT NULL
    )
);

END DO

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

Запись в таблицу

Запрос читает события из топика и записывает их в таблицу output_table. Таблица должна быть создана заранее со схемой, соответствующей выбираемым колонкам.

Важно

Запись в таблицы в потоковых запросах поддерживается только в режиме UPSERT. Операция INSERT INTO не поддерживается, так как при повторной обработке событий (гарантия at-least-once) она привела бы к дублированию строк. При UPSERT, если строка с таким первичным ключом уже существует, она будет обновлена, иначе будет вставлена новая строка, a INSERT INTO завершится с ошибкой.

CREATE STREAMING QUERY write_table_example AS
DO BEGIN

-- Запись в таблицу (только UPSERT, INSERT не поддерживается)
UPSERT INTO output_table
SELECT
    Id,
    Name
FROM
    ext_source.input_topic -- или локальный топик input_topic
WITH (
    FORMAT = json_each_row,  -- Формат входных данных
    SCHEMA = (               -- Схема входных данных
        Id Uint64 NOT NULL,
        Name Utf8 NOT NULL
    )
);

END DO

Подробнее: Запись в таблицы.

См. также