Чтение данных из Data Streams с помощью соединений в Query
Соединения удобно использовать для прототипирования и первоначальной настройки подключения к данным Yandex Data Streams.
Yandex Data Streams — сервис для передачи потоков данных нескольким приложениям. Каждое приложение обрабатывает данные независимо от других.
Пример чтения данных в формате JSON из Yandex Data Streams:
SELECT
JSON_VALUE(CAST(Data AS Json), "$.action") AS action
FROM yds.`input_stream`
WITH (
format=raw,
SCHEMA
(
Data String
)
)
LIMIT 10;
Примечание
Данные из потокового источника передаются в виде бесконечного потока. Чтобы остановить обработку и получить результат в консоли, данные в примере ограничены с помощь оператора LIMIT, который задает количество строк результата.
Настройка соединения
Чтобы настроить чтение данных из Yandex Data Streams:
-
Перейдите
в сервис Yandex Query. -
На панели слева выберите Соединения.
-
Нажмите Создать.
-
В открывшемся окне в поле Имя укажите название соединения с Yandex Data Streams.
-
В выпадающем поле Тип выберите
Data Streams. -
В поле Облако и каталог выберите расположение источника данных.
-
В выпадающем поле База данных выберите базу данных Yandex Managed Service for YDB, где ранее был создан поток Yandex Data Streams.
-
В поле Сервисный аккаунт выберите сервисный аккаунт, который будет использоваться для чтения данных, или создайте новый, выдав ему права
yds.editor.Чтобы использовать сервисный аккаунт, пользователю нужна роль
iam.serviceAccounts.user. -
Нажмите Создать.
Модель данных
Данные через Yandex Data Streams передаются в бинарном виде. Для чтения данных используйте SQL-выражение следующего вида:
SELECT
<выражение>
FROM
<соединение>.<имя_потока>
WITH
(
format=raw,
SCHEMA
(
Data String
)
)
WHERE
<фильтр>;
Где:
<выражение>— выражение, определяющее результат запроса;<соединение>— название соединения с потоком данных Data Streams, созданного в предыдущем разделе;<имя_потока>— название потока данных в Data Streams;<фильтр>— условие фильтрации данных.
Пример чтения данных
Пример запроса для чтения данных из Yandex Data Streams:
$data =
SELECT
JSON_VALUE(Data, "$.host") AS host,
JSON_VALUE(Data, "$.count") 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";
SELECT
*
FROM
$data
LIMIT 10;
Где:
|
Поле |
Тип |
Описание |
|
|
Название соединения с Yandex Data Streams. |
|
|
|
Название потока-источника данных. |
|
|
|
Строка |
Название хоста. |
|
|
Строка |
Количество событий. |
|
|
Строка |
Тег события. |
|
|
Строка |
Формат данных. На данный момент поддерживается только формат |