Запись данных из Yandex Query в потоки Yandex Data Streams
Yandex Data Streams — сервис для передачи потоков данных нескольким приложениям. Каждое приложение обрабатывает данные независимо от других.
Пример записи данных в формате JSON в Yandex Data Streams:
INSERT INTO yds.`output_stream`
SELECT
ToBytes(Unwrap(Json::SerializeJson(Yson::From(
<|"predefined":
<|
"host": host,
"count": count,
|>,
"optional":
<|
"tag": tag
|>
|>))))
FROM
$data;
Настройка соединения
Чтобы настроить запись данных в Yandex Data Streams:
- Перейдите
в сервис Yandex Query. - На панели слева выберите Соединения.
- Нажмите Создать.
- В открывшемся окне в поле Имя укажите название соединения с Yandex Data Streams.
- В поле Тип выберите
Data Streams. - В поле База данных выберите базу данных Yandex Managed Service for YDB, где ранее был создан поток Yandex Data Streams.
- В поле Сервисный аккаунт выберите сервисный аккаунт, который будет использоваться для записи данных, или создайте новый и назначьте ему роль
yds.writer. - Нажмите Создать.
Модель данных
Данные через Yandex Data Streams передаются в бинарном виде. Запись данных выполняется с помощью SQL-выражений и в общем случае выглядит следующим образом:
INSERT INTO <соединение>.<имя_потока>
<выражение>
FROM
<запрос>
Где:
<соединение>— название соединения с потоком данных Data Streams, созданного в предыдущем разделе.<имя_потока>— название потока данных в Data Streams.<выражение>— выражение, определяющее записываемые данные.<запрос>— запрос-источник данных Yandex Query.
Пример записи данных
Пример запроса для чтения данных из Yandex Data Streams и записи результатов в Yandex Data Streams:
$data =
SELECT
JSON_VALUE(Data, "$.host") AS host,
CAST(JSON_VALUE(Data, "$.count") AS Int) AS count,
JSON_VALUE(Data, "$.tag") AS tag,
FROM
(
SELECT
CAST(Data AS Json) AS Data
FROM yds.`input_stream`
WITH(
format=raw,
SCHEMA
(
Data String
)
)
)
WHERE
JSON_VALUE(Data, "$.tag") = "my_tag";
INSERT INTO yds.`output_stream`
SELECT
ToBytes(Unwrap(Json::SerializeJson(Yson::From(
<|"predefined":
<|
"host": host,
"count": count,
|>,
"optional":
<|
"tag": tag
|>
|>))))
FROM
$data;
Где:
|
Поле |
Тип |
Описание |
|
|
Название соединения с Yandex Data Streams. |
|
|
|
Название потока — источника данных в SQL-запросе. |
|
|
|
Название потока — приемника данных в SQL-запросе. |
|
|
|
Строка |
Строковый параметр запроса. |
|
|
Целое число |
Числовой параметр запроса. |
|
|
Строка |
Формат данных. На данный момент поддерживается только формат |
Результаты обработки записываются в выходной поток Yandex Data Streams. Чтобы упростить обработку, результаты преобразуются в формат JSON с помощью следующей конструкции:
ToBytes(Unwrap(Json::SerializeJson(Yson::From(
<|"key": value|>,
<|"key2":
<|"child_key": child_value|>,
|>,
))))
В документации YQL приведено подробное описание модулей Yson
Поддерживаемые форматы записи
В Data Streams можно выполнять запись только в виде байтового потока, который интепретируется на принимающей стороне.
Настройки форматов файлов и алгоритмов сжатия при записи в Data Streams не применяются.