---
metadata:
  - name: generator
    content: Diplodoc Platform v5.52.0
alternate:
  - https://ydb.tech/docs/ru/dev/streaming-query/S3-enrichment.md?version=v25.4
  - href: ru/dev/streaming-query/S3-enrichment.md
    type: text/markdown
    title: Markdown version
  - href: ../../llms.txt
    type: text/markdown
    title: llms.txt
sourcePath: ru/core/dev/streaming-query/S3-enrichment.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ydb.tech/docs/ru/llms.txt

# Обогащение данных (S3)

Обогащение данных — добавление к событиям из потока дополнительной информации из справочника. Например, событие содержит только идентификатор, а справочник позволяет добавить к нему название или другие атрибуты.

В [потоковых запросах](https://ydb.tech/docs/ru/concepts/streaming-query.md?version=v25.4) можно присоединить к потоку данные, хранимые в S3, с помощью конструкции `JOIN`. Поток должен быть слева, справочник из S3 — справа.


{% note warning %}

Справочник полностью загружается в память при запуске запроса. Если данные в S3 изменились, для получения актуальной версии справочника необходимо перезапустить запрос — удалить его с помощью [DROP STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/drop-streaming-query.md?version=v25.4) и создать заново с помощью [CREATE STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/create-streaming-query.md?version=v25.4).

{% endnote %}

Справочник хранится в S3 и подключается через [внешний источник данных](https://ydb.tech/docs/ru/concepts/query_execution/federated_query/s3/external_data_source.md?version=v25.4).

## Подготовка источников данных

Создаём два внешних источника: один для YDB (топики), другой для S3 (справочник). Для хранения токена используется [секрет](https://ydb.tech/docs/ru/yql/reference/syntax/create-secret.md?version=v25.4), источники создаются через [CREATE EXTERNAL DATA SOURCE](https://ydb.tech/docs/ru/yql/reference/syntax/create-external-data-source.md?version=v25.4).

```sql
-- Секрет с токеном для подключения к YDB
CREATE SECRET `secrets/ydb_token` WITH (value = "<ydb_token>");

-- Источник данных YDB для чтения/записи топиков
CREATE EXTERNAL DATA SOURCE ydb_source WITH (
    SOURCE_TYPE = "Ydb",
    LOCATION = "<ydb_endpoint>",
    DATABASE_NAME = "<db_name>",
    AUTH_METHOD = "TOKEN",
    TOKEN_SECRET_NAME = "secrets/ydb_token"
);

-- Источник данных S3 для чтения справочника
CREATE EXTERNAL DATA SOURCE s3_source WITH (
    SOURCE_TYPE = "ObjectStorage",
    LOCATION = "<s3_endpoint>",
    AUTH_METHOD = "NONE"
)
```

Где:

- `<ydb_endpoint>` — эндпоинт YDB, например `grpcs://<ydb_host>:2135` для Yandex Cloud.
- `<db_name>` — путь к базе данных YDB, например `/Root/database`.
- `<s3_endpoint>` — URL S3-хранилища, например `https://storage.yandexcloud.net/<bucket>` для Yandex Cloud.

## Создание потокового запроса

Запрос читает события из входного топика, присоединяет к каждому событию название сервиса из справочника по `ServiceId` и записывает результат в выходной топик. Запрос создаётся с помощью [CREATE STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/create-streaming-query.md?version=v25.4).

Подробнее о форматах данных (`json_each_row`, `csv_with_names` и др.): [Форматы данных при чтении/записи из топиков](https://ydb.tech/docs/ru/dev/streaming-query/streaming-query-formats.md?version=v25.4).

```sql
CREATE STREAMING QUERY query_with_join AS
DO BEGIN

-- Чтение событий из входного топика
$topic_data = SELECT
    *
FROM
    ydb_source.input_topic
WITH (
    FORMAT = json_each_row,
    SCHEMA = (
        Time String NOT NULL,
        ServiceId Uint32 NOT NULL,
        Message String NOT NULL
    )
);

-- Чтение справочника сервисов из S3
$s3_data = SELECT
    *
FROM
    s3_source.`file.csv`
WITH (
    FORMAT = csv_with_names,
    SCHEMA = (
        ServiceId Uint32,
        Name Utf8
    )
);

-- Присоединение справочника к потоку по ServiceId
$joined_data = SELECT
    s.Name AS Name,
    t.*
FROM
    $topic_data AS t
LEFT JOIN
    $s3_data AS s
ON
    t.ServiceId = s.ServiceId;

-- Запись результата в выходной топик в формате JSON
INSERT INTO
    ydb_source.output_topic
SELECT
    ToBytes(Unwrap(Yson::SerializeJson(Yson::From(TableRow()))))
FROM
    $joined_data;

END DO
```

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

- [TableRow](https://ydb.tech/docs/ru/yql/reference/builtins/basic.md?version=v25.4#tablerow)
- [Yson::From](https://ydb.tech/docs/ru/yql/reference/udf/list/yson.md?version=v25.4#ysonfrom)
- [Yson::SerializeJson](https://ydb.tech/docs/ru/yql/reference/udf/list/yson.md?version=v25.4#ysonserializejson)
- [Unwrap](https://ydb.tech/docs/ru/yql/reference/builtins/basic.md?version=v25.4#unwrap)
- [ToBytes](https://ydb.tech/docs/ru/yql/reference/builtins/basic.md?version=v25.4#to-from-bytes).
