---
metadata:
  - name: generator
    content: Diplodoc Platform v5.50.6
alternate:
  - https://ydb.tech/docs/en/reference/ydb-sdk/topic.md
  - https://ydb.tech/docs/ru/reference/ydb-sdk/topic.md
sourcePath: ru/core/reference/ydb-sdk/topic.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ydb.tech/docs/ru/llms.txt

# Работа с топиками

В этой статье приведены примеры использования YDB SDK для работы с [топиками](https://ydb.tech/docs/ru/concepts/datamodel/topic.md).

Перед выполнением примеров [создайте топик](https://ydb.tech/docs/ru/reference/ydb-cli/topic-create.md) и [добавьте читателя](https://ydb.tech/docs/ru/reference/ydb-cli/topic-consumer-add.md).

## Примеры работы с топиками

{% list tabs group=lang %}

- C++

  [Пример читателя на GitHub](https://github.com/ydb-platform/ydb/tree/main/ydb/public/sdk/cpp/examples/topic_reader)

- Go

  [Примеры на GitHub](https://github.com/ydb-platform/ydb-go-sdk/tree/master/examples/topic)

- Java

  [Примеры на GitHub](https://github.com/ydb-platform/ydb-java-examples/tree/master/ydb-cookbook/src/main/java/tech/ydb/examples/topic)

- Python

  [Примеры на GitHub](https://github.com/ydb-platform/ydb-python-sdk/tree/main/examples/topic)

- C#

  [Примеры на GitHub](https://github.com/ydb-platform/ydb-dotnet-sdk/tree/main/examples/src/Topic)

- JavaScript

  [Примеры на GitHub](https://github.com/ydb-platform/ydb-js-sdk/tree/main/examples/topic)

- Rust

  [Примеры на GitHub](https://github.com/ydb-platform/ydb-rs-sdk/tree/master/ydb/examples) (`topic-writer`, `topic-reader-retry`, `topic-read-in-transaction-example`).

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

## Инициализация соединения с топиками {#init}

{% list tabs group=lang %}

- C++

  Для работы с топиками создаются экземпляры драйвера YDB и клиента.

  Драйвер YDB отвечает за взаимодействие приложения и YDB на транспортном уровне. Драйвер должен существовать на всем протяжении жизненного цикла работы с топиками и должен быть инициализирован перед созданием клиента.

  Клиент сервиса топиков ([исходный код](https://github.com/ydb-platform/ydb/blob/d2d07d368cd8ffd9458cc2e33798ee4ac86c733c/ydb/public/sdk/cpp/client/ydb_topic/topic.h#L1589)) работает поверх драйвера YDB и отвечает за управляющие операции с топиками, а также создание сессий чтения и записи.

  Фрагмент кода приложения для инициализации драйвера YDB:

  ```cpp
  // Create driver instance.
  auto driverConfig = NYdb::TDriverConfig()
      .SetEndpoint(opts.Endpoint)
      .SetDatabase(opts.Database)
      .SetAuthToken(std::getenv("YDB_TOKEN"));

  NYdb::TDriver driver(driverConfig);
  ```

  В этом примере используется аутентификационный токен, сохранённый в переменной окружения `YDB_TOKEN`. Подробнее про [соединение с БД](https://ydb.tech/docs/ru/concepts/connect.md) и [аутентификацию](https://ydb.tech/docs/ru/security/authentication.md).

  Фрагмент кода приложения для создания клиента:

  ```cpp
  NYdb::NTopic::TTopicClient topicClient(driver);
  ```

- Go

  Для работы с топиками используется экземпляр драйвера YDB, созданный с помощью `ydb.Open`. Клиент топиков доступен через метод `db.Topic()`.

  ```go
  package main

  import (
    "context"
    "os"

    "github.com/ydb-platform/ydb-go-sdk/v3"
    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topicoptions"
  )

  func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    db, err := ydb.Open(ctx,
      os.Getenv("YDB_CONNECTION_STRING"),
    )
    if err != nil {
      panic(err)
    }
    defer db.Close(ctx)

    // db.Topic() — клиент для работы с топиками
    writer, err := db.Topic().StartWriter("topic-path")
    if err != nil {
      panic(err)
    }

    reader, err := db.Topic().StartReader("consumer-name",
      topicoptions.ReadTopic("topic-path"),
    )
    if err != nil {
      panic(err)
    }
    _ = writer
    _ = reader
  }
  ```

- Java

  Для работы с топиками создаются экземпляры транспорта YDB и клиента.

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

  Фрагмент кода приложения для инициализации транспорта YDB:

  ```java
  try (GrpcTransport transport = GrpcTransport.forConnectionString(connString)
          .withAuthProvider(CloudAuthHelper.getAuthProviderFromEnviron())
          .build()) {
      // Use YDB transport
  }
  ```

  В этом примере используется вспомогательный метод `CloudAuthHelper.getAuthProviderFromEnviron()`, получающий токен из переменных окружения.
  Например, `YDB_ACCESS_TOKEN_CREDENTIALS`.
  Подробнее про [соединение с БД](https://ydb.tech/docs/ru/concepts/connect.md) и [аутентификацию](https://ydb.tech/docs/ru/security/authentication.md).

  Клиент сервиса топиков ([исходный код](https://github.com/ydb-platform/ydb-java-sdk/blob/master/topic/src/main/java/tech/ydb/topic/TopicClient.java#L34)) работает поверх транспорта YDB и отвечает как за управляющие операции с топиками, так и за создание писателей и читателей.

  Фрагмент кода приложения для создания клиента:

  ```java
  try (TopicClient topicClient = TopicClient.newClient(transport)
                .setCompressionExecutor(compressionExecutor)
                .build()) {
    // Use topic client
  }
  ```

  В обоих примерах кода выше используется блок ([try-with-resources](https://docs.oracle.com/javase/tutorial/essential/exceptions/tryResourceClose.html)).
  Это позволяет автоматически закрывать клиент и транспорт при выходе из этого блока, т.к. оба являются наследниками `AutoCloseable`.

- C#

  Для работы с топиками достаточно передать строку подключения напрямую в конструктор нужного клиента.

  В этом примере используется анонимная аутентификация. Подробнее про [соединение с базой данных](https://ydb.tech/docs/ru/concepts/connect.md) и [аутентификацию](https://ydb.tech/docs/ru/security/authentication.md).

  Фрагмент кода приложения для создания различных клиентов к топикам:

  ```c#
  const string connectionString = "Host=localhost;Port=2136;Database=/local";

  await using var topicClient = new TopicClient(connectionString);

  await using var writer = new WriterBuilder<string>(connectionString, topicName)
  {
      ProducerId = "ProducerId_Example"
  }.Build();

  await using var reader = new ReaderBuilder<string>(connectionString)
  {
      ConsumerName = "Consumer_Example",
      SubscribeSettings = { new SubscribeSettings(topicName) }
  }.Build();
  ```

- Python

  Для работы с топиками создаётся экземпляр драйвера YDB. Клиент топиков доступен через атрибут `topic_client` и используется для управляющих операций с топиками, а также создания писателей и читателей.

  {% list tabs %}
  - Native SDK

    ```python
    import os
    import ydb

    driver_config = ydb.DriverConfig(
        endpoint=os.environ["YDB_ENDPOINT"],
        database=os.environ["YDB_DATABASE"],
    )
    driver = ydb.Driver(driver_config)
    driver.wait(timeout=5)
    # driver.topic_client — клиент для работы с топиками
    writer = driver.topic_client.writer(topic_path)
    reader = driver.topic_client.reader(topic=topic_path, consumer=consumer_name)
    ```

  - Native SDK (Asyncio)

    ```python
    import os
    import ydb

    driver_config = ydb.DriverConfig(
        endpoint=os.environ["YDB_ENDPOINT"],
        database=os.environ["YDB_DATABASE"],
    )
    async with ydb.aio.Driver(driver_config) as driver:
        await driver.wait(timeout=5)
        # driver.topic_client — клиент для работы с топиками
        writer = driver.topic_client.writer(topic_path)
        reader = driver.topic_client.reader(topic=topic_path, consumer=consumer_name)
    ```

  {% endlist %}

  Подробнее про [соединение с БД](https://ydb.tech/docs/ru/concepts/connect.md) и [аутентификацию](https://ydb.tech/docs/ru/security/authentication.md).

- JavaScript

  ```javascript
  const t = topic(driver);

  await using reader = t.createReader({
    topic: "/Root/demo-topic",
    consumer: "demo-consumer",
  });

  await using writer = t.createWriter({
    topic: "/Root/demo-topic",
    producer: "demo-producer",
  });
  ```

- Rust

  ```rust
  use ydb::{ClientBuilder, YdbResult};

  #[tokio::main]
  async fn main() -> YdbResult<()> {
      let client = ClientBuilder::new_from_connection_string("grpc://localhost:2136/local")?.client()?;
      client.wait().await?;
      let mut topic_client = client.topic_client();
      // topic_client.create_reader(...), create_writer_with_params(...), ...
      Ok(())
  }
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

## Управление топиками {#manage}

### Создание топика {#create-topic}

Единственный обязательный параметр для создания топика - это его путь, остальные параметры опциональны.

{% list tabs group=lang %}

- C++

  Полный список настроек можно посмотреть [в заголовочном файле](https://github.com/ydb-platform/ydb/blob/d2d07d368cd8ffd9458cc2e33798ee4ac86c733c/ydb/public/sdk/cpp/client/ydb_topic/topic.h#L394).

  Пример создания топика c тремя партициями и поддержкой кодека ZSTD:

  ```cpp
  auto settings = NYdb::NTopic::TCreateTopicSettings()
      .PartitioningSettings(3, 3)
      .AppendSupportedCodecs(NYdb::NTopic::ECodec::ZSTD);

  auto status = topicClient
      .CreateTopic("my-topic", settings)  // returns TFuture<TStatus>
      .GetValueSync();
  ```

- Go

  Полный список поддерживаемых параметров можно посмотреть в [документации SDK](https://pkg.go.dev/github.com/ydb-platform/ydb-go-sdk/v3/topic/topicoptions#CreateOption).

  Пример создания топика со списком поддерживаемых кодеков и минимальным количеством партиций

  ```go
  err := db.Topic().Create(ctx, "topic-path",
    // optional
    topicoptions.CreateWithSupportedCodecs(topictypes.CodecRaw, topictypes.CodecGzip),

    // optional
    topicoptions.CreateWithMinActivePartitions(3),
  )
  ```

- Python

  Пример создания топика со списком поддерживаемых кодеков и минимальным количеством партиций

  {% list tabs %}
  - Native SDK

    ```python
    driver.topic_client.create_topic(topic_path,
        supported_codecs=[ydb.TopicCodec.RAW, ydb.TopicCodec.GZIP], # optional
        min_active_partitions=3,                                    # optional
    )
    ```

  - Native SDK (Asyncio)

    ```python
    await driver.topic_client.create_topic(topic_path,
        supported_codecs=[ydb.TopicCodec.RAW, ydb.TopicCodec.GZIP],  # optional
        min_active_partitions=3,                                     # optional
    )
    ```

  {% endlist %}

- Java

  Полный список настроек можно посмотреть [в коде SDK](https://github.com/ydb-platform/ydb-java-sdk/blob/master/topic/src/main/java/tech/ydb/topic/settings/CreateTopicSettings.java#L97).

  Пример создания топика со списком поддерживаемых кодеков и минимальным количеством партиций

  ```java
  topicClient.createTopic(topicPath, CreateTopicSettings.newBuilder()
                  // Optional
                  .setSupportedCodecs(SupportedCodecs.newBuilder()
                          .addCodec(Codec.RAW)
                          .addCodec(Codec.GZIP)
                          .build())
                  // Optional
                  .setPartitioningSettings(PartitioningSettings.newBuilder()
                          .setMinActivePartitions(3)
                          .build())
                  .build());
  ```

- C#

  Пример создания топика со списком поддерживаемых кодеков и минимальным количеством партиций:

  ```c#
  await topicClient.CreateTopic(new CreateTopicSettings
  {
      Path = topicName,
      Consumers = { new Consumer("Consumer_Example") },
      SupportedCodecs = { Codec.Raw, Codec.Gzip },
      PartitioningSettings = new PartitioningSettings
      {
          MinActivePartitions = 3
      }
  });
  ```

- JavaScript

  ```javascript
  const topicService = driver.createClient(TopicServiceDefinition);
  await topicService.createTopic(
    create(CreateTopicRequestSchema, {
      path: "/path-to-my-topic",
      partitioningSettings: {
        minActivePartitions: 1n,
        maxActivePartitions: 100n,
      },
      consumers: [{ name: "my-consumer" }],
    }),
  );
  ```

- Rust

  ```rust
  use ydb::{Codec, CreateTopicOptionsBuilder, YdbResult};

  topic_client
      .create_topic(
          "/local/my-topic".into(),
          CreateTopicOptionsBuilder::default()
              .min_active_partitions(3)
              .supported_codecs(vec![Codec::Raw, Codec::Zstd])
              .build()?,
      )
      .await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Изменение топика {#alter-topic}

{% list tabs group=lang %}

- C++

  При изменении топика в параметрах метода `AlterTopic` нужно указать путь топика и параметры, которые будут изменяться. Изменяемые параметры представлены структурой `TAlterTopicSettings`.

  Полный список настроек можно посмотреть [в заголовочном файле](https://github.com/ydb-platform/ydb/blob/d2d07d368cd8ffd9458cc2e33798ee4ac86c733c/ydb/public/sdk/cpp/client/ydb_topic/topic.h#L458).

  Пример добавления [важного читателя](../../concepts/datamodel/topic#important-consumer) к топику и установки [времени хранения сообщений](../../concepts/datamodel/topic#retention-time) для топика в два дня:

  ```cpp
  auto alterSettings = NYdb::NTopic::TAlterTopicSettings()
      .BeginAddConsumer("my-consumer")
          .Important(true)
      .EndAddConsumer()
      .SetRetentionPeriod(TDuration::Days(2));

  auto status = topicClient
      .AlterTopic("my-topic", alterSettings)  // returns TFuture<TStatus>
      .GetValueSync();
  ```

- Go

  При изменении топика в параметрах нужно указать путь топика и те параметры, которые будут изменяться.

  Полный список поддерживаемых параметров можно посмотреть в [документации SDK](https://pkg.go.dev/github.com/ydb-platform/ydb-go-sdk/v3/topic/topicoptions#AlterOption).

  Пример добавления читателя к топику

  ```go
  err := db.Topic().Alter(ctx, "topic-path",
    topicoptions.AlterWithAddConsumers(topictypes.Consumer{
      Name:            "new-consumer",
      SupportedCodecs: []topictypes.Codec{topictypes.CodecRaw, topictypes.CodecGzip}, // optional
    }),
  )
  ```

- Python

  Пример изменения списка поддерживаемых кодеков и минимального количества партиций у топика

  {% list tabs %}
  - Native SDK

    ```python
    driver.topic_client.alter_topic(topic_path,
        set_supported_codecs=[ydb.TopicCodec.RAW, ydb.TopicCodec.GZIP], # optional
        set_min_active_partitions=3,                                    # optional
    )
    ```

  - Native SDK (Asyncio)

    ```python
    await driver.topic_client.alter_topic(topic_path,
        set_supported_codecs=[ydb.TopicCodec.RAW, ydb.TopicCodec.GZIP],  # optional
        set_min_active_partitions=3,                                     # optional
    )
    ```

  {% endlist %}

- Java

  При изменении топика в параметрах метода `alterTopic` нужно указать путь топика и параметры, которые будут изменяться.

  Полный список настроек можно посмотреть [в коде SDK](https://github.com/ydb-platform/ydb-java-sdk/blob/master/topic/src/main/java/tech/ydb/topic/settings/AlterTopicSettings.java#L23).

  ```java
  topicClient.alterTopic(topicPath, AlterTopicSettings.newBuilder()
                  .addAddConsumer(Consumer.newBuilder()
                          .setName("new-consumer")
                          .setSupportedCodecs(SupportedCodecs.newBuilder()
                                  .addCodec(Codec.RAW)
                                  .addCodec(Codec.GZIP)
                                  .build())
                          .build())
                  .build());


- JavaScript

  ```javascript
  const topicService = driver.createClient(TopicServiceDefinition);
  await topicService.alterTopic(
    create(AlterTopicRequestSchema, {
      path: "/path-to-my-topic",
      addConsumers: [{ name: "my-consumer-2" }],
    }),
  );
  ```

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- JavaScript

  ```javascript
  const topicService = driver.createClient(TopicServiceDefinition);
  await topicService.alterTopic(
    create(AlterTopicRequestSchema, {
      path: "/path-to-my-topic",
      addConsumers: [{ name: "my-consumer-2" }],
    }),
  );
  ```

- Rust

  ```rust
  use ydb::{AlterTopicOptionsBuilder, YdbResult};

  topic_client
      .alter_topic(
          "/local/my-topic".into(),
          AlterTopicOptionsBuilder::default()
              .set_min_active_partitions(Some(5))
              .build()?,
      )
      .await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Получение информации о топике {#describe-topic}

{% list tabs group=lang %}

- C++

  Для получения информации о топике используется метод `DescribeTopic`.

  Описание топика представлено структурой `TTopicDescription`.

  Полный список полей описания смотри [в заголовочном файле](https://github.com/ydb-platform/ydb/blob/d2d07d368cd8ffd9458cc2e33798ee4ac86c733c/ydb/public/sdk/cpp/client/ydb_topic/topic.h#L163).

  Получить доступ к этому описанию можно так:

  ```cpp
  auto result = topicClient.DescribeTopic("my-topic").GetValueSync();
  if (result.IsSuccess()) {
      const auto& description = result.GetTopicDescription();
      std::cout << "Topic description: " << GetProto(description) << std::endl;
  }
  ```

  Существует отдельный метод для получения информации о читателе - `DescribeConsumer`.

- Go

  ```go
    descResult, err := db.Topic().Describe(ctx, "topic-path")
  if err != nil {
    log.Fatalf("failed describe topic: %v", err)
    return
  }
  fmt.Printf("describe: %#v\n", descResult)
  ```

- Python

  {% list tabs %}
  - Native SDK

    ```python
    info = driver.topic_client.describe_topic(topic_path)
    print(info)
    ```

  - Native SDK (Asyncio)

    ```python
    info = await driver.topic_client.describe_topic(topic_path)
    print(info)
    ```

  {% endlist %}

- Java

  Для получения информации о топике используется метод `describeTopic`.

  Полный список полей описания можно посмотреть [в коде SDK](https://github.com/ydb-platform/ydb-java-sdk/blob/master/topic/src/main/java/tech/ydb/topic/description/TopicDescription.java#L19).

  ```java
  Result<TopicDescription> topicDescriptionResult = topicClient.describeTopic(topicPath)
          .join();
  TopicDescription description = topicDescriptionResult.getValue();
  ```

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- JavaScript

  ```javascript
  const topicService = driver.createClient(TopicServiceDefinition);
  await topicService.describeTopic(
    create(DescribeTopicRequestSchema, {
      path: "/path-to-my-topic",
    }),
  );
  ```

- Rust

  ```rust
  use ydb::{DescribeTopicOptionsBuilder, YdbResult};

  let description = topic_client
      .describe_topic(
          "/local/my-topic".into(),
          DescribeTopicOptionsBuilder::default().include_stats(true).build()?,
      )
      .await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Удаление топика {#drop-topic}

Для удаления топика достаточно указать путь к нему.

{% list tabs group=lang %}

- C++

  ```cpp
  auto status = topicClient.DropTopic("my-topic").GetValueSync();
  ```

- Go

  ```go
    err := db.Topic().Drop(ctx, "topic-path")
  ```

- Python

  {% list tabs %}
  - Native SDK

    ```python
    driver.topic_client.drop_topic(topic_path)
    ```

  - Native SDK (Asyncio)

    ```python
    await driver.topic_client.drop_topic(topic_path)
    ```

  {% endlist %}

- Java

  ```java
  topicClient.dropTopic(topicPath);
  ```

- C#

  ```c#
  await topicClient.DropTopic(topicName);
  ```

- JavaScript

  ```javascript
  const topicService = driver.createClient(TopicServiceDefinition);
  await topicService.dropTopic(
    create(DropTopicRequestSchema, {
      path: "/path-to-my-topic",
    }),
  );
  ```

- Rust

  ```rust
  topic_client.drop_topic("/local/my-topic".into()).await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

## Запись сообщений {#write}

### Подключение к топику для записи сообщений {#start-writer}

На данный момент поддерживается подключение только с совпадающими идентификаторами [источника и группы сообщений](../../concepts/datamodel/topic#producer-id) (`producer_id` и `message_group_id`), в будущем это ограничение будет снято.

{% list tabs group=lang %}

- C++

  Подключение к топику на запись представлено объектом сессии записи с интерфейсом `IWriteSession` или `ISimpleBlockingWriteSession` (вариант для простой записи по одному сообщению без подтверждения, блокирующейся при превышении числа inflight записей или размера буфера SDK). Настройки сессии записи представлены структурой `TWriteSessionSettings`, для варианта `ISimpleBlockingWriteSession` часть настроек не поддерживается.

  Полный список настроек смотри [в заголовочном файле](https://github.com/ydb-platform/ydb/blob/d2d07d368cd8ffd9458cc2e33798ee4ac86c733c/ydb/public/sdk/cpp/client/ydb_topic/topic.h#L1199).

  Пример создания сессии записи с интерфейсом `IWriteSession`.

  ```cpp
  std::string producerAndGroupID = "group-id";
  auto settings = NYdb::NTopic::TWriteSessionSettings()
      .Path("my-topic")
      .ProducerId(producerAndGroupID)
      .MessageGroupId(producerAndGroupID);

  auto session = topicClient.CreateWriteSession(settings);
  ```

- Go

  ```go
  producerAndGroupID := "group-id"
  writer, err := db.Topic().StartWriter(producerAndGroupID, "topicName",
    topicoptions.WithMessageGroupID(producerAndGroupID),
  )
  if err != nil {
      return err
  }
  ```

- Python

  {% list tabs %}
  - Native SDK

    ```python
    writer = driver.topic_client.writer(topic_path)
    ```

  - Native SDK (Asyncio)

    ```python
    writer = driver.topic_client.writer(topic_path)
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    Инициализация настроек писателя:

    ```java
    String producerAndGroupID = "group-id";
    WriterSettings settings = WriterSettings.newBuilder()
          .setTopicPath(topicPath)
          .setProducerId(producerAndGroupID)
          .setMessageGroupId(producerAndGroupID)
          .build();
    ```

    Создание синхронного писателя:

    ```java
    SyncWriter writer = topicClient.createSyncWriter(settings);
    ```

    После создания писателя его необходимо инициализировать. Для этого есть два метода:

    - `init()`: неблокирующий, запускает процесс инициализации в фоне и не ждёт его завершения.

      ```java
      writer.init();
      ```

    - `initAndWait()`: блокирующий, запускает процесс инициализации и ждёт его завершения. Если в процессе инициализации возникла ошибка, будет брошено исключение.

      ```java
      try {
          writer.initAndWait();
          logger.info("Init finished successfully");
      } catch (Exception exception) {
          logger.error("Exception while initializing writer: ", exception);
          return;
      }
      ```

  - Асинхронный API

    Инициализация настроек писателя:

    ```java
    String producerAndGroupID = "group-id";
    WriterSettings settings = WriterSettings.newBuilder()
          .setTopicPath(topicPath)
          .setProducerId(producerAndGroupID)
          .setMessageGroupId(producerAndGroupID)
          .build();
    ```

    Создание и инициализация асинхронного писателя:

    ```java
    AsyncWriter writer = topicClient.createAsyncWriter(settings);

    // Init in background
    writer.init()
            .thenRun(() -> logger.info("Init finished successfully"))
            .exceptionally(ex -> {
                logger.error("Init failed with ex: ", ex);
                return null;
            });
    ```

    {% endlist %}

- C#

  ```c#
  await using var writer = new WriterBuilder<string>(connectionString, topicName)
  {
      ProducerId = "ProducerId_Example"
  }.Build();
  ```

- JavaScript

  ```javascript
  await using writer = createTopicWriter(driver, {
    topic: topicName,
    producer: producerName,
  });
  ```

- Rust

  ```rust
  use ydb::{TopicWriter, TopicWriterOptionsBuilder, YdbResult};

  let writer: TopicWriter = topic_client
      .create_writer_with_params(
          TopicWriterOptionsBuilder::default()
              .topic_path("/local/my-topic".into())
              .producer_id("group-id".into())
              .message_group_id("group-id".into())
              .build()?,
      )
      .await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Запись сообщений {#writing-messages}

{% list tabs group=lang %}

- C++

  Асинхронная запись возможна через интерфейс `IWriteSession`.

  Работа пользователя с объектом `IWriteSession` в общем устроена как обработка цикла событий с тремя типами событий: `TReadyToAcceptEvent`, `TAcksEvent` и `TSessionClosedEvent`.

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

  Если обработчик для некоторого события не установлен, его необходимо получить и обработать в методах `GetEvent` / `GetEvents`. Для неблокирующего ожидания очередного события есть метод `WaitEvent` с интерфейсом `TFuture<void>()`.

  Для записи каждого сообщения пользователь должен "потратить" move-only объект `TContinuationToken`, который выдаёт SDK с событием `TReadyToAcceptEvent`. При записи сообщения можно установить пользовательские seqNo и временную метку создания, но по умолчанию их проставляет SDK автоматически.

  По умолчанию `Write` выполняется асинхронно - данные из сообщений вычитываются и сохраняются во внутренний буфер, отправка происходит в фоне в соответствии с настройками `MaxMemoryUsage`, `MaxInflightCount`, `BatchFlushInterval`, `BatchFlushSizeBytes`. Сессия сама переподключается к YDB при обрывах связи и повторяет отправку сообщений пока это возможно, в соответствии с настройкой `RetryPolicy`. При получении ошибки, после которой невозможно продолжить работу, сессия чтения отправляет пользователю `TSessionClosedEvent` с диагностической информацией.

  Так может выглядеть запись нескольких сообщений в цикле событий без использования обработчиков:

  ```cpp
  // Event loop
  while (true) {
      // Get event
      // May block for a while if write session is busy
      std::optional<NYdb::NTopic::TWriteSessionEvent::TEvent> event = session->GetEvent(/*block=*/true);

      if (auto* readyEvent = std::get_if<NYdb::NTopic::TWriteSessionEvent::TReadyToAcceptEvent>(&*event)) {
          session->Write(std::move(event.ContinuationToken), "This is yet another message.");

      } else if (auto* ackEvent = std::get_if<NYdb::NTopic::TWriteSessionEvent::TAcksEvent>(&*event)) {
          std::cout << ackEvent->DebugString() << std::endl;

      } else if (auto* closeSessionEvent = std::get_if<NYdb::NTopic::TSessionClosedEvent>(&*event)) {
          break;
      }
  }
  ```

- Go

  Для отправки сообщения - достаточно в поле Data сохранить Reader, из которого можно будет прочитать данные. Можно рассчитывать на то что данные каждого сообщения читаются один раз (или до первой ошибки), к моменту возврата из Write данные будут уже прочитаны и сохранены во внутренний буфер.

  SeqNo и дата создания сообщений по умолчанию проставляются автоматически.

  По умолчанию Write выполняется асинхронно - данные из сообщений вычитываются и сохраняются во внутренний буфер, отправка происходит в фоне. Writer сам переподключается к YDB при обрывах связи и повторяет отправку сообщений пока это возможно. При получении ошибки, после которой невозможно продолжить работу, Writer останавливается и следующие вызовы Write будут завершаться с ошибкой.

  ```go
  err := writer.Write(ctx,
    topicwriter.Message{Data: strings.NewReader("1")},
    topicwriter.Message{Data: bytes.NewReader([]byte{1,2,3})},
    topicwriter.Message{Data: strings.NewReader("3")},
  )
  if err != nil {
    return err
  }
  ```

- Python

  Для отправки сообщений можно передавать как просто содержимое сообщения (bytes, str), так и вручную задавать некоторые свойства. Объекты можно передавать по одному или сразу в массиве (list). Метод `write` выполняется асинхронно. Возврат из метода происходит сразу после того как сообщения будут положены во внутренний буфер клиента, обычно это происходит быстро. Ожидание может возникнуть, если внутренний буфер уже заполнен и нужно подождать, пока часть данных будет отправлена на сервер.

  {% list tabs %}
  - Native SDK

    ```python
    # Простая отправка сообщений, без явного указания метаданных.
    # Удобно начинать, удобно использовать пока важно только содержимое сообщения.
    writer = driver.topic_client.writer(topic_path)
    writer.write("mess")  # Строки будут переданы в кодировке utf-8, так удобно отправлять
                          # текстовые сообщения.
    writer.write(bytes([1, 2, 3]))  # Эти байты будут отправлены "как есть", так удобно отправлять
                                    # бинарные данные.
    writer.write(["mess-1", "mess-2"])  # Здесь за один вызов отправляется несколько сообщений —
                                         # так снижаются накладные расходы на внутренние процессы SDK,
                                         # имеет смысл при большом потоке сообщений.

    # Полная форма, используется, когда кроме содержимого сообщения нужно вручную задать и его свойства.
    writer = driver.topic_client.writer(topic="topic-path", auto_seqno=False, auto_created_at=False)

    writer.write(ydb.TopicWriterMessage("asd", seqno=123, created_at=datetime.datetime.now()))
    writer.write(ydb.TopicWriterMessage(bytes([1, 2, 3]), seqno=124, created_at=datetime.datetime.now()))

    # В полной форме так же можно отправлять несколько сообщений за один вызов функции.
    # Это имеет смысл при большом потоке отправляемых сообщений — для снижения
    # накладных расходов на внутренние вызовы SDK.
    writer.write([
      ydb.TopicWriterMessage("asd", seqno=123, created_at=datetime.datetime.now()),
      ydb.TopicWriterMessage(bytes([1, 2, 3]), seqno=124, created_at=datetime.datetime.now(),
      ])
    ```

  - Native SDK (Asyncio)

    ```python
    writer = driver.topic_client.writer(topic_path)
    await writer.write("mess")
    await writer.write(bytes([1, 2, 3]))
    await writer.write(["mess-1", "mess-2"])
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    Метод `send` блокирует управление, пока сообщение не будет помещено в очередь отправки.
    Попадание сообщения в эту очередь означает, что писатель сделает всё возможное для доставки сообщения.
    Например, если сессия записи по какой-то причине оборвётся, писатель переустановит соединение и попробует отправить это сообщение на новой сессии.
    Но попадание сообщения в очередь отправки не гарантирует того, что сообщение в итоге будет записано.
    Например, могут возникать ошибки, приводящие к завершению работы писателя до того, как сообщения из очереди будут отправлены.
    Если нужно подтверждение успешной записи для каждого сообщения, используйте асинхронного писателя и проверяйте статус, возвращаемый методом `send`.

    ```java
    writer.send(Message.of("11".getBytes()));

    long timeoutSeconds = 5; // How long should we wait for a message to be put into sending buffer
    try {
      writer.send(
              Message.newBuilder()
                      .setData("22".getBytes())
                      .setCreateTimestamp(Instant.now().minusSeconds(5))
                      .build(),
              timeoutSeconds,
              TimeUnit.SECONDS
      );
    } catch (TimeoutException exception) {
      logger.error("Send queue is full. Couldn't put message into sending queue within {} seconds", timeoutSeconds);
    } catch (InterruptedException | ExecutionException exception) {
      logger.error("Couldn't put the message into sending queue due to exception: ", exception);
    }
    ```

  - Асинхронный API

    Метод `send` в асинхронном клиенте неблокирующий. Помещает сообщение в очередь отправки.
    Метод возвращает `CompletableFuture<WriteAck>`, позволяющую проверить, действительно ли сообщение было записано.
    В случае, если очередь переполнена, будет брошено исключение QueueOverflowException.
    Это способ сигнализировать пользователю о том, что поток записи следует притормозить.
    В таком случае стоит или пропускать сообщения, или выполнять повторные попытки записи через exponential backoff.
    Также можно увеличить размер клиентского буфера (`setMaxSendBufferMemorySize`), чтобы обрабатывать больший объем сообщений перед тем, как он заполнится.

    ```java
    try {
      // Non-blocking. Throws QueueOverflowException if send queue is full
      writer.send(Message.of("33".getBytes()));
    } catch (QueueOverflowException exception) {
      // Send queue is full. Need to retry with backoff or skip
    }
    ```

  {% endlist %}

- C#

  Асинхронная запись сообщения в топик.

  ```c#
  var asyncWriteTask = writer.WriteAsync("Hello, Example YDB Topics!"); // Task<WriteResult>
  ```

- JavaScript

  ```javascript
  // Пишет сообщение во внутренний буфер
  writer.write(Buffer.from("Hello, world!", "utf-8"));

  // Для немедленной отправки нужно вызвать flush
  await writer.flush();

  // Или закрыть писатель
  await writer.close();
  ```

- Rust

  ```rust
  use ydb::{TopicWriterMessageBuilder, YdbResult};

  writer
      .write(
          TopicWriterMessageBuilder::default()
              .data(b"payload".to_vec())
              .build()?,
      )
      .await?;
  writer.stop().await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Запись сообщений с подтверждением о сохранении на сервере

{% list tabs group=lang %}

- C++

  Получение подтверждений от сервера возможно через интерфейс `IWriteSession`.

  Ответы о записи сообщений на сервере приходят клиенту SDK в виде событий `TAcksEvent`. В одном событии могут содержаться ответы о нескольких отправленных ранее сообщениях. Варианты ответа: запись подтверждена (`EES_WRITTEN`), запись отброшена как дубликат ранее записанного сообщения (`EES_ALREADY_WRITTEN`) или запись отброшена по причине сбоя (`EES_DISCARDED`).

  Пример установки обработчика TAcksEvent для сессии записи:

  ```cpp
  auto settings = NYdb::NTopic::TWriteSessionSettings()
    // other settings are set here
    .EventHandlers(
      NYdb::NTopic::TWriteSessionSettings::TEventHandlers()
        .AcksHandler(
          [&](NYdb::NTopic::TWriteSessionEvent::TAcksEvent& event) {
            for (const auto& ack : event.Acks) {
              if (ack.State == NYdb::NTopic::TWriteSessionEvent::TWriteAck::EEventState::EES_WRITTEN) {
                ackedSeqNo.insert(ack.SeqNo);
                std::cout << "Acknowledged message with seqNo " << ack.SeqNo << std::endl;
              }
            }
          }
        )
    );

  auto session = topicClient.CreateWriteSession(settings);
  ```

  В такой сессии записи события `TAcksEvent` не будут приходить пользователю в `GetEvent` / `GetEvents`, вместо этого SDK при получении подтверждений от сервера будет вызывать переданный обработчик. Аналогично можно настраивать обработчики на остальные типы событий.

- Go

  При подключении можно указать опцию синхронной записи сообщений - topicoptions.WithSyncWrite(true). Тогда Write будет возвращаться только после того как получит подтверждение с сервера о сохранении всех, сообщений переданных в вызове. При этом SDK так же как и обычно будет при необходимости переподключаться и повторять отправку сообщений. В этом режиме контекст управляет только временем ожидания ответа из SDK, т.е. даже после отмены контекста SDK продолжит попытки отправить сообщения.

  ```go
  producerAndGroupID := "group-id"
  writer, _ := db.Topic().StartWriter(producerAndGroupID, "topicName",
    topicoptions.WithMessageGroupID(producerAndGroupID),
    topicoptions.WithSyncWrite(true),
  )

  err = writer.Write(ctx,
    topicwriter.Message{Data: strings.NewReader("1")},
    topicwriter.Message{Data: bytes.NewReader([]byte{1,2,3})},
    topicwriter.Message{Data: strings.NewReader("3")},
  )
  if err != nil {
    return err
  }
  ```

- Python

  Есть два способа получить подтверждение о записи сообщений на сервере:
  - `flush()` — дожидается подтверждения для всех сообщений, записанных ранее во внутренний буфер.
  - `write_with_ack(...)` — отправляет сообщение и ждет подтверждение его доставки от сервера. При отправке нескольких сообщений подряд это способ работает медленно.

  {% list tabs %}
  - Native SDK

    ```python
    # Положить несколько сообщений во внутренний буфер, затем дождаться,
    # пока все они будут доставлены до сервера.
    for mess in messages:
        writer.write(mess)

    writer.flush()

    # Можно отправить несколько сообщений и дождаться подтверждения на всю группу.
    writer.write_with_ack(["mess-1", "mess-2"])

    # Ожидание при отправке каждого сообщения — этот метод вернет результат только после получения
    # подтверждения от сервера.
    # Это самый медленный вариант отправки сообщений, используйте его только если такой режим
    # действительно нужен.
    writer.write_with_ack("message")
    ```

  - Native SDK (Asyncio)

    ```python
    for mess in messages:
        await writer.write(mess)

    await writer.flush()

    await writer.write_with_ack(["mess-1", "mess-2"])
    await writer.write_with_ack("message")
    ```

  {% endlist %}

- Java

  Метод `send` возвращает `CompletableFuture<WriteAck>`. Её успешное завершение означает подтверждение записи сервером.
  В структуре `WriteAck` содержится информация о seqNo, offset и статусе записи:

  ```java
  writer.send(Message.of(message))
        .whenComplete((result, ex) -> {
            if (ex != null) {
                logger.error("Exception on writing message message: ", ex);
            } else {
                switch (result.getState()) {
                    case WRITTEN:
                        WriteAck.Details details = result.getDetails();
                        StringBuilder str = new StringBuilder("Message was written successfully");
                        if (details != null) {
                            str.append(", offset: ").append(details.getOffset());
                        }
                        logger.debug(str.toString());
                        break;
                    case ALREADY_WRITTEN:
                        logger.warn("Message has already been written");
                        break;
                    default:
                        break;
                }
            }
        });
  ```

- C#

  Асинхронная запись сообщения в топик. В случае переполнения внутреннего буфера будет ожидать, когда буфер освободится для повторной отправки.

  ```c#
  await writer.WriteAsync("Hello, Example YDB Topics!");
  ```

  В случае, если сервер недоступен, сообщения могут накапливаться в очереди в ожидании отправки. Для управления временем ожидания можно использовать токен отмены (`CancellationToken`). Однако, при таком подходе существует вероятность того, что пользователь может отменить отправку уже подтвержденного сообщения.

  ```c#
  var writeCts = new CancellationTokenSource();
  writeCts.CancelAfter(TimeSpan.FromSeconds(3));

  await writer.WriteAsync("Hello, Example YDB Topics!", writeCts.Token);
  ```

- JavaScript

  Все сообщения записываются во внутренний буфер. Для отправки на сервер есть 3 механизма: два автоматических и один ручной. Ручной - это вызов метода `writer.flush` который возвращает последний seqno записанный на сервере. Автоматическая отправка происходит по условиям:
  - Превышение размера внутреннего буфера `maxBufferBytes` (значение по умолчанию = 256MiB).
  - По тику интервала периодической отправки `flushIntervalMs` (значение по умолчанию = 10ms).

  ```javascript
  await using writer = createTopicWriter(driver, {
    topic: topicName,
    producer: producerName,
    // Callback that is called when writer receives an acknowledgment for a message.
    onAck: (seqNo, status) => {
      console.log("ACK", seqNo, status);
    },
  })

  writer.write(Buffer.from("Hello, world!", "utf-8"));

  // Чтобы получить последний записанный seqNo на сервере.
  await writer.flush();
  ```

- Rust

  ```rust
  use ydb::{TopicWriterMessageBuilder, YdbResult};

  writer
      .write_with_ack(
          TopicWriterMessageBuilder::default()
              .data(b"payload".to_vec())
              .build()?,
      )
      .await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Выбор кодека для сжатия сообщений {#codec}

Подробнее о [сжатии данных в топиках](../../concepts/datamodel/topic#message-codec).

{% list tabs group=lang %}

- C++

  Сжатие, которое используется при отправке сообщений методом `Write`, задаётся при [создании сессии записи](#start-writer) настройками `Codec` и `CompressionLevel`. По умолчанию выбирается кодек GZIP.
  Пример создания сессии записи без сжатия сообщений:

  ```cpp
  auto settings = NYdb::NTopic::TWriteSessionSettings()
    // other settings are set here
    .Codec(ECodec::RAW);

  auto session = topicClient.CreateWriteSession(settings);
  ```

  Если необходимо в рамках сессии записи отправить сообщение, сжатое другим кодеком, можно использовать метод `WriteEncoded` с указанием кодека и размера расжатого сообщения. Для успешной записи этим способом используемый кодек должен быть разрешён в настройках топика.

- Go

  По умолчанию SDK выбирает кодек автоматически (с учетом настроек топика). В автоматическом режиме SDK сначала отправляет по одной группе сообщений каждым из разрешенных кодеков, затем иногда будет пробовать сжать сообщения всеми доступными кодеками и выбирать кодек, дающий наименьший размер сообщения. Если для топика список разрешенных кодеков пуст, то автовыбор производится между Raw и Gzip-кодеками.

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

  ```go
  producerAndGroupID := "group-id"
  writer, _ := db.Topic().StartWriter(producerAndGroupID, "topicName",
    topicoptions.WithMessageGroupID(producerAndGroupID),
    topicoptions.WithCodec(topictypes.CodecGzip),
  )
  ```

- Python

  По умолчанию SDK выбирает кодек автоматически (с учетом настроек топика). В автоматическом режиме SDK сначала отправляет по одной группе сообщений каждым из разрешенных кодеков, затем иногда будет пробовать сжать сообщения всеми доступными кодеками и выбирать кодек, дающий наименьший размер сообщения. Если для топика список разрешенных кодеков пуст, то автовыбор производится между Raw и Gzip-кодеками.

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

  ```python
  writer = driver.topic_client.writer(topic_path,
      codec=ydb.TopicCodec.GZIP,
  )
  ```

- Java

  ```java
  String producerAndGroupID = "group-id";
  WriterSettings settings = WriterSettings.newBuilder()
          .setTopicPath(topicPath)
          .setProducerId(producerAndGroupID)
          .setMessageGroupId(producerAndGroupID)
          .setCodec(Codec.ZSTD)
          .build();
  ```

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- JavaScript

  ```javascript
  await using writer = t.createWriter({
    codec: Codec.RAW,
  });

  await using writer = t.createWriter({
    codec: Codec.GZIP,
  });

  await using writer = t.createWriter({
    codec: Codec.LZOP,
  });

  await using writer = t.createWriter({
    codec: 10000, // CUSTOM (допустимый диапазон: 10000–19999)
  });
  ```

- Rust

  Выбор кодека сжатия при записи в Rust SDK пока недоступен; сообщения отправляются с кодеком `Raw`.

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#341](https://github.com/ydb-platform/ydb-rs-sdk/issues/341)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Запись сообщений без дедупликации {#nodedup}

Подробнее о записи без дедупликации — в [соответствующем разделе концепций](../../concepts/datamodel/topic#no-dedup).

{% list tabs group=lang %}

- C++

  Если в настройках сессии записи не указывается опция `ProducerId`, будет создана сессия записи без дедупликации.
  Пример создания такой сессии записи:

  ```cpp
  auto settings = NYdb::NTopic::TWriteSessionSettings()
      .Path(myTopicPath);

  auto session = topicClient.CreateWriteSession(settings);
  ```

  Для включения дедупликации нужно в настройках сессии записи указать опцию `ProducerId` или явно включить дедупликацию, вызвав метод `DeduplicationEnabled()`, например, как в секции ["Подключение к топику"](#start-writer).

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Go

  В **ydb-go-sdk** при создании писателя, если не передавать `topicoptions.WithWriterProducerID`, SDK всё равно подставляет идентификатор производителя (генерирует его автоматически). Режим записи без дедупликации, эквивалентный отсутствию `ProducerId` в примере для C++ выше, в текущей версии SDK недоступен.

- Java

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#341](https://github.com/ydb-platform/ydb-rs-sdk/issues/341)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Запись метаданных на уровне сообщения {#messagemeta}

При записи сообщения можно дополнительно указать метаданные как список пар "ключ-значение". Эти данные будут доступны при вычитывании сообщения.
Ограничение на размер метаданных — не более 1000 ключей.

{% list tabs group=lang %}

- C++

  Воспользоваться функцией записи метаданных можно с помощью метода `Write()`, принимающего `TWriteMessage` объект:

  ```cpp
  auto settings = NYdb::NTopic::TWriteSessionSettings()
      .Path(myTopicPath)
  // set all other settings;
  ;

  auto session = topicClient.CreateWriteSession(settings);

  std::optional<NYdb::NTopic::TWriteSessionEvent::TEvent> event = session->GetEvent(/*block=*/true);
  NYdb::NTopic::TWriteMessage message("This is yet another message").MessageMeta({
      {"meta-key", "meta-value"},
      {"another-key", "value"}
  });

  if (auto* readyEvent = std::get_if<NYdb::NTopic::TWriteSessionEvent::TReadyToAcceptEvent>(&*event)) {
      session->Write(std::move(event.ContinuationToken), std::move(message));
  }
  ```

- Go

  Метаданные задаются в поле `Metadata` структуры `topicwriter.Message`:

  ```go
  err := writer.Write(ctx, topicwriter.Message{
    Data: strings.NewReader("message-data"),
    Metadata: map[string][]byte{
      "meta-key":    []byte("meta-value"),
      "another-key": []byte("value"),
    },
  })
  ```

  При чтении метаданные доступны в поле `Metadata` у сообщения:

  ```go
  msg, err := reader.ReadMessage(ctx)
  if err != nil {
    return err
  }
  for k, v := range msg.Metadata {
    fmt.Printf("%s: %s\n", k, string(v))
  }
  ```

- Java

  При конструировании сообщения для записи с помощью Builder'а, ему можно передать объекты типа `MetadataItem` с парой ключ типа `String` + значение типа `byte[]`.

  Можно передать сразу `List` таких объектов:

  ```java
  List<MetadataItem> metadataItems = Arrays.asList(
          new MetadataItem("meta-key", "meta-value".getBytes()),
          new MetadataItem("another-key", "value".getBytes())
  );
  writer.send(
          Message.newBuilder()
                  .setMetadataItems(metadataItems)
                  .build()
  );
  ```

  Или добавлять каждый `MetadataItem` отдельно:

  ```java
  writer.send(
          Message.newBuilder()
                  .addMetadataItem(new MetadataItem("meta-key", "meta-value".getBytes()))
                  .addMetadataItem(new MetadataItem("another-key", "value".getBytes()))
                  .build()
  );
  ```

  При чтении эти метаданные сообщения получить, вызвав на нём метод `getMetadataItems()`:

  ```java
  Message message = reader.receive();
  List<MetadataItem> metadata = message.getMetadataItems();
  ```

- Python

  Для использования функции передачи метаданных создайте объект `TopicWriterMessage` с аргументом `metadata_items`, как показано ниже:

  {% list tabs %}
  - Native SDK

    ```python
    message = ydb.TopicWriterMessage(data=f"message-data", metadata_items={"meta-key": "meta-value"})
    writer.write(message)
    ```

  - Native SDK (Asyncio)

    ```python
    message = ydb.TopicWriterMessage(data="message-data", metadata_items={"meta-key": "meta-value"})
    await writer.write(message)
    ```

  {% endlist %}

  Во время чтения метаданные можно получить из поля `metadata_items` объекта `PublicMessage`:

  ```python
  message = reader.receive_message()
  for meta_key, meta_value in message.metadata_items.items():
      print(f"{meta_key}: {meta_value}")
  ```

- C#

  ```c#
  await writer.WriteAsync(
      new Ydb.Sdk.Services.Topic.Writer.Message<string>("Hello, Example YDB Topics!")
          { Metadata = { new Metadata("meta-key", "meta-value"u8.ToArray()) } }
  );
  ```

- JavaScript

  ```javascript
  writer.write(Buffer.from("Hello, world!", "utf-8"), {
    metadataItems: {
      "meta-key": new TextEncoder().encode("meta-value"),
    },
  });
  ```

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#341](https://github.com/ydb-platform/ydb-rs-sdk/issues/341)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Запись в транзакции {#write-tx}

{% list tabs group=lang %}

- C++

  Для записи в топик в транзакции необходимо передать ссылку на объект транзакции в метод `Write` сессии записи.

  [Пример на GitHub](https://github.com/ydb-platform/ydb-cpp-sdk/blob/main/examples/topic_writer/transaction/main.cpp)

  ```c++
  NYdb::NQuery::TQueryClient queryClient(driver);

  NYdb::NStatusHelpers::ThrowOnError(queryClient.RetryQuerySync([](NYdb::NQuery::TSession session) -> NYdb::TStatus {
      auto beginTxResult = session.BeginTransaction().GetValueSync();
      if (!beginTxResult.IsSuccess()) {
          return beginTxResult;
      }
      auto tx = beginTxResult.GetTransaction();

      NYdb::NTopic::TWriteMessage writeMessage("message");

      topicSession->Write(std::move(writeMessage), tx);
      return tx.Commit().GetValueSync();
  }));
  ```

- Go

  Для записи в топик в транзакции необходимо создать транзакционного писателя через вызов [TopicClient.StartTransactionalWriter](https://pkg.go.dev/github.com/ydb-platform/ydb-go-sdk/v3/topic#Client.StartTransactionalWriter). После этого можно отправлять сообщения, как обычно. Закрывать транзакционного писателя не требуется — это происходит автоматически при завершении транзакции.

  [Пример на GitHub](https://github.com/ydb-platform/ydb-go-sdk/blob/master/examples/topic/topicwriter/topic_writer_transaction.go)

  ```go
  err := db.Query().DoTx(ctx, func(ctx context.Context, tx query.TxActor) error {
    writer, err := db.Topic().StartTransactionalWriter(tx, topicName)
    if err != nil {
      return err
    }

    return writer.Write(ctx, topicwriter.Message{Data: strings.NewReader("asd")})
  })
  ```

- Python

  Для записи в топик в транзакции необходимо создать транзакционного писателя через вызов `topic_client.tx_writer`. После этого можно отправлять сообщения, как обычно. Закрывать транзакционного писателя не требуется — это происходит автоматически при завершении транзакции.

  В примере ниже нет явного вызова `tx.commit()` — он происходит неявно при успешном завершении лямбды `callee`.

  [Пример на GitHub](https://github.com/ydb-platform/ydb-python-sdk/blob/main/examples/topic/topic_transactions_example.py)

  {% list tabs %}
  - Native SDK

    ```python
    with ydb.QuerySessionPool(driver) as session_pool:

        def callee(tx: ydb.QueryTxContext):
            tx_writer: ydb.TopicTxWriter = driver.topic_client.tx_writer(tx, topic)

            for i in range(message_count):
                result_stream = tx.execute(query=f"select {i} as res;")
                for result_set in result_stream:
                    message = str(result_set.rows[0]["res"])
                    tx_writer.write(ydb.TopicWriterMessage(message))
                    print(f"Message {message} was written with tx.")

        session_pool.retry_tx_sync(callee)
    ```

  - Native SDK (Asyncio)

    [Пример на GitHub](https://github.com/ydb-platform/ydb-python-sdk/blob/main/examples/topic/topic_transactions_async_example.py)

    ```python
    async with ydb.aio.QuerySessionPool(driver) as session_pool:

        async def callee(tx: ydb.aio.QueryTxContext):
            tx_writer: ydb.TopicTxWriterAsyncIO = driver.topic_client.tx_writer(tx, topic)

            for i in range(message_count):
                async with await tx.execute(query=f"select {i} as res;") as result_stream:
                    async for result_set in result_stream:
                        message = str(result_set.rows[0]["res"])
                        await tx_writer.write(ydb.TopicWriterMessage(message))
                        print(f"Message {result_set.rows[0]['res']} was written with tx.")

        await session_pool.retry_tx_async(callee)
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    [Пример на GitHub](https://github.com/ydb-platform/ydb-java-examples/blob/develop/ydb-cookbook/src/main/java/tech/ydb/examples/topic/transactions/TransactionWriteSync.java)

    В настройках `SendSettings` метода `send` можно указать транзакцию.
    Тогда сообщение будет записано вместе с коммитом этой транзакцией.

    ```java
    // creating a session in the table service
    Result<Session> sessionResult = tableClient.createSession(Duration.ofSeconds(10)).join();
    if (!sessionResult.isSuccess()) {
      logger.error("Couldn't get a session from the pool: {}", sessionResult);
      return; // retry or shutdown
    }
    Session session = sessionResult.getValue();
    // creating a transaction in the table service
    // this transaction is not yet active and has no id
    TableTransaction transaction = session.createNewTransaction(TxMode.SERIALIZABLE_RW);

    // get message text within the transaction
    Result<DataQueryResult> dataQueryResult = transaction.executeDataQuery("SELECT \"Hello, world!\";")
          .join();
    if (!dataQueryResult.isSuccess()) {
      logger.error("Couldn't execute DataQuery: {}", dataQueryResult);
      return; // retry or shutdown
    }
    // now the transaction is active and has an id

    ResultSetReader rsReader = dataQueryResult.getValue().getResultSet(0);
    byte[] message;
    if (rsReader.next()) {
      message = rsReader.getColumn(0).getBytes();
    } else {
      return; // retry or shutdown
    }

    writer.send(
          Message.of(message),
          SendSettings.newBuilder()
                  .setTransaction(transaction)
                  .build()
    );

    // flush to wait until all messages reach server before commit
    writer.flush();

    Status commitStatus = transaction.commit().join();
    analyzeCommitStatus(commitStatus);
    ```

  - Асинхронный API

    [Пример на GitHub](https://github.com/ydb-platform/ydb-java-examples/blob/develop/ydb-cookbook/src/main/java/tech/ydb/examples/topic/transactions/TransactionWriteAsync.java)

    В настройках `SendSettings` метода `send` можно указать транзакцию.
    Тогда сообщение будет записано вместе с коммитом этой транзакцией.

    ```java
    // creating a session in the table service
    Result<Session> sessionResult = tableClient.createSession(Duration.ofSeconds(10)).join();
    if (!sessionResult.isSuccess()) {
      logger.error("Couldn't get a session from the pool: {}", sessionResult);
      return; // retry or shutdown
    }
    Session session = sessionResult.getValue();
    // creating a transaction in the table service
    // this transaction is not yet active and has no id
    TableTransaction transaction = session.createNewTransaction(TxMode.SERIALIZABLE_RW);

    // get message text within the transaction
    Result<DataQueryResult> dataQueryResult = transaction.executeDataQuery("SELECT \"Hello, world!\";")
          .join();
    if (!dataQueryResult.isSuccess()) {
      logger.error("Couldn't execute DataQuery: {}", dataQueryResult);
      return; // retry or shutdown
    }
    // now the transaction is active and has an id

    ResultSetReader rsReader = dataQueryResult.getValue().getResultSet(0);
    byte[] message;
    if (rsReader.next()) {
      message = rsReader.getColumn(0).getBytes();
    } else {
      return; // retry or shutdown
    }

    try {
      writer.send(Message.newBuilder()
                              .setData(message)
                              .build(),
                      SendSettings.newBuilder()
                              .setTransaction(transaction)
                              .build())
              .whenComplete((result, ex) -> {
                  if (ex != null) {
                      logger.error("Exception while sending a message: ", ex);
                  } else {
                      switch (result.getState()) {
                          case WRITTEN:
                              WriteAck.Details details = result.getDetails();
                              logger.info("Message was written successfully, offset: " + details.getOffset());
                              break;
                          case ALREADY_WRITTEN:
                              logger.info("Message has already been written");
                              break;
                          default:
                              break;
                      }
                  }
              })
              // Waiting for the message to reach the server before committing the transaction
              .join();

      Status commitStatus = transaction.commit().join();
      analyzeCommitStatus(commitStatus);
    } catch (QueueOverflowException exception) {
      logger.error("Queue overflow exception while sending a message{}: ", index, exception);
      // Send queue is full. Need to retry with backoff or skip
    }
    ```

  {% endlist %}

  <!-- source: ru/reference/ydb-sdk/_includes/alerts/java_transaction_requirements.md -->
  {% note info %}

  Требования к транзакции:

  * Это должна быть активная (имеющая идентификатор) транзакция в одном из сервисов YDB. Например, [Table](https://github.com/ydb-platform/ydb-java-sdk/blob/master/table/src/main/java/tech/ydb/table/transaction/TableTransaction.java) или [Query](https://github.com/ydb-platform/ydb-java-sdk/blob/master/query/src/main/java/tech/ydb/query/QueryTransaction.java).
  * В Topic Service поддерживается только уровень изоляции транзакций `SERIALIZABLE_RW`.

  {% endnote %}
  <!-- endsource: ru/reference/ydb-sdk/_includes/alerts/java_transaction_requirements.md -->

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#341](https://github.com/ydb-platform/ydb-rs-sdk/issues/341)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

## Чтение сообщений {#reading}

### Подключение к топику для чтения сообщений {#start-reader}

Чтение сообщений из топика может выполнятся с указанием Consumer'а, связанного с этим топиком, а также без Consumer'а. Если Consumer не указан, то клиентское приложение должно самостоятельно рассчитывать offset для чтения сообщений. Более подробно пример чтения без Consumer'а рассмотрен в [соответствующей секции](#no-consumer).

Создать Consumer можно при [создании](#create-topic) или [изменении](#alter-topic) топика.
У топика может быть несколько Consumer'ов и для каждого из них сервер хранит свой прогресс чтения.

{% list tabs group=lang %}

- C++

  Подключение для чтения из одного или нескольких топиков представлено объектом сессии чтения с интерфейсом `IReadSession`. Настройки сессии чтения представлены структурой `TReadSessionSettings`.

  Полный список настроек смотри [в заголовочном файле](https://github.com/ydb-platform/ydb/blob/d2d07d368cd8ffd9458cc2e33798ee4ac86c733c/ydb/public/sdk/cpp/client/ydb_topic/topic.h#L1344).

  Чтобы создать подключение к существующему топику `my-topic` через добавленного ранее читателя `my-consumer`, используйте следующий код:

  ```cpp
  auto settings = NYdb::NTopic::TReadSessionSettings()
      .ConsumerName("my-consumer")
      .AppendTopics("my-topic");

  auto session = topicClient.CreateReadSession(settings);
  ```

- Go

  Чтобы создать подключение к существующему топику `my-topic` через добавленного ранее читателя `my-consumer`, используйте следующий код:

  ```go
  reader, err := db.Topic().StartReader("my-consumer", topicoptions.ReadTopic("my-topic"))
  if err != nil {
      return err
  }
  ```

- Python

  Чтобы создать подключение к существующему топику `my-topic` через добавленного ранее читателя `my-consumer`, используйте следующий код:

  {% list tabs %}
  - Native SDK

    ```python
    reader = driver.topic_client.reader(topic="my-topic", consumer="my-consumer")
    ```

  - Native SDK (Asyncio)

    ```python
    reader = driver.topic_client.reader(topic="my-topic", consumer="my-consumer")
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    Инициализация настроек читателя

    ```java
    ReaderSettings settings = ReaderSettings.newBuilder()
          .setConsumerName(consumerName)  // имя consumer'а, зарегистрированного на топике
          .addTopic(TopicReadSettings.newBuilder()
                  .setPath(topicPath)
                  .setReadFrom(Instant.now().minus(Duration.ofHours(24))) // читать с этой временной метки (опционально)
                  .setMaxLag(Duration.ofMinutes(30)) // максимальное отставание от конца очереди (опционально)
                  .build())
          .build();
    ```

    Создание синхронного читателя

    ```java
    SyncReader reader = topicClient.createSyncReader(settings);
    ```

    После создания синхронного читателя необходимо инициализировать. Для этого следует воспользоваться одним их двух методов:
    - `init()`: неблокирующий, запускает процесс инициализации в фоне и не ждёт его завершения.

    ```java
    reader.init();
    ```

    - `initAndWait()`: блокирующий, запускает процесс инициализации и ждёт его завершения. Если в процессе инициализации возникла ошибка, будет брошено исключение.

    ```java
    try {
        reader.initAndWait();
        logger.info("Init finished successfully");
    } catch (Exception exception) {
        logger.error("Exception while initializing reader: ", exception);
        return;
    }
    ```

  - Асинхронный API

    Инициализация настроек читателя

    ```java
    ReaderSettings settings = ReaderSettings.newBuilder()
          .setConsumerName(consumerName)  // имя consumer'а, зарегистрированного на топике
          .addTopic(TopicReadSettings.newBuilder()
                  .setPath(topicPath)
                  .setReadFrom(Instant.now().minus(Duration.ofHours(24))) // читать с этой временной метки (опционально)
                  .setMaxLag(Duration.ofMinutes(30)) // максимальное отставание от конца очереди (опционально)
                  .build())
          .build();
    ```

    Для асинхронного читателя, помимо общих настроек чтения `ReaderSettings`, понадобятся настройки обработчика событий `ReadEventHandlersSettings`, в которых необходимо передать экземпляр наследника `ReadEventHandler`.
    Он будет описывать, как должна происходить обработка различных событий, происходящих во время чтения.

    ```java
    ReadEventHandlersSettings handlerSettings = ReadEventHandlersSettings.newBuilder()
          .setEventHandler(new Handler())
          .build();
    ```

    Опционально, в `ReadEventHandlersSettings` можно указать executor'а, на котором будет происходить обработка сообщений; по умолчанию используется внутренний поток SDK.

    Для реализации обработчика событий можно унаследоваться от `AbstractReadEventHandler` и переопределить метод `onMessages`.
    Метод `onMessages` вызывается каждый раз, когда SDK получает очередной пакет сообщений от сервера. В рамках одного вызова приходит один или несколько сообщений, которые можно подтвердить (`commit`) как по отдельности, так и после обработки всего пакета. Пример реализации:

    ```java
    private class Handler extends AbstractReadEventHandler {
      @Override
      public void onMessages(DataReceivedEvent event) {
          for (Message message : event.getMessages()) {
              StringBuilder str = new StringBuilder();
              logger.info("Message received. SeqNo={}, offset={}", message.getSeqNo(), message.getOffset());

              process(message);

              message.commit().thenRun(() -> {
                  logger.info("Message committed");
              });
          }
      }
    }
    ```

    Создание и инициализация асинхронного читателя:

    ```java
    AsyncReader reader = topicClient.createAsyncReader(readerSettings, handlerSettings);
    // Init in background
    reader.init()
          .thenRun(() -> logger.info("Init finished successfully"))
          .exceptionally(ex -> {
              logger.error("Init failed with ex: ", ex);
              return null;
          });
    ```

  {% endlist %}

- C#

  ```c#
  await using var reader = new ReaderBuilder<string>(connectionString)
  {
      ConsumerName = "Consumer_Example",
      SubscribeSettings = { new SubscribeSettings(topicName) }
  }.Build();
  ```

- JavaScript

  ```javascript
  await using reader = createTopicReader(driver, {
    topic: topicName,
    consumer: consumerName,
  });
  ```

- Rust

  ```rust
  use ydb::YdbResult;

  let mut reader = topic_client
      .create_reader("my-consumer", "/local/my-topic")
      .await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

Вы также можете использовать расширенный вариант создания подключения, чтобы указать несколько топиков и задать параметры чтения. Следующий код создаст подключение к топикам `my-topic` и `my-specific-topic` через читателя `my-consumer`:

{% list tabs group=lang %}

- C++

  ```cpp
  auto settings = NYdb::NTopic::TReadSessionSettings()
      .ConsumerName("my-consumer")
      .AppendTopics("my-topic")
      .AppendTopics(
          NYdb::NTopic::TTopicReadSettings("my-specific-topic")
              .ReadFromTimestamp(someTimestamp)
      );

  auto session = topicClient.CreateReadSession(settings);
  ```

- Go

  ```go
  reader, err := db.Topic().StartReader("my-consumer", []topicoptions.ReadSelector{
      {
          Path: "my-topic",
      },
      {
          Path:       "my-specific-topic",
          ReadFrom:   time.Date(2022, 7, 1, 10, 15, 0, 0, time.UTC),
      },
      },
  )
  if err != nil {
      return err
  }
  ```

  Также в примере выше задаётся время, с которого следует начинать читать сообщения.

- Python

  Функциональность находится в разработке.

- Java

  ```java
  ReaderSettings settings = ReaderSettings.newBuilder()
          .setConsumerName(consumerName)
          .addTopic(TopicReadSettings.newBuilder()
                  .setPath("my-topic")
                  .build())
          .addTopic(TopicReadSettings.newBuilder()
                  .setPath("my-specific-topic")
                  .setReadFrom(Instant.now().minus(Duration.ofHours(24))) // Optional
                  .setMaxLag(Duration.ofMinutes(30)) // Optional
                  .build())
          .build();
  ```

- C#

  ```c#
  await using var reader = new ReaderBuilder<string>(connectionString)
  {
      ConsumerName = "Consumer_Example",
      SubscribeSettings =
      {
          new SubscribeSettings(topicName),
          new SubscribeSettings(topicName + "_another") { ReadFrom = DateTime.Now }
      }
  }.Build();
  ```

- JavaScript

  ```javascript
  await using reader = createTopicReader(driver, {
    topic: {
      path: topicPath,
      partitionIds: [1n, 2n, 3n],
    },
    consumer: consumerName,
  });

  await using reader = createTopicReader(driver, {
    topic: {
      path: topicPath,
      maxLag: "1s", // number, import('ms').StringValue, protobuff Duration
    },
    consumer: consumerName,
  });

  await using reader = createTopicReader(driver, {
    topic: {
      path: topicPath,
      readFrom: new Date(), // number, Date, protobuf Timestamp
    },
    consumer: consumerName,
  });

  await using reader = createTopicReader(driver, {
    topic: [
      {
        path: topicPath,
        partitionIds: [1n, 2n, 3n],
      },
      {
        path: topicPath2,
        maxLag: "1s",
      },
      {
        path: topicPath3,
        readFrom: new Date(),
      },
      // ...
    ],
    consumer: consumerName,
  });
  ```

- Rust

  ```rust
  use std::time::{Duration, SystemTime};

  use ydb::{TopicReaderOptionsBuilder, TopicSelector, TopicSelectors, YdbResult};

  let mut reader = topic_client
      .create_reader_with_params(
          TopicReaderOptionsBuilder::default()
              .consumer("my-consumer".into())
              .topic(TopicSelectors(vec![
                  TopicSelector::new("/local/my-topic"),
                  TopicSelector {
                      path: "/local/my-specific-topic".into(),
                      partition_ids: None,
                      read_from: Some(SystemTime::now() - Duration::from_secs(3600)),
                  },
              ]))
              .build()?,
      )
      .await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Чтение сообщений {#reading-messages}

Сервер хранит [позицию чтения сообщений](https://ydb.tech/docs/ru/concepts/datamodel/topic.md#consumer-offset). После вычитывания очередного сообщения клиент может [отправить на сервер подтверждение обработки](#commit). Позиция чтения изменится, а при новом подключении будут вычитаны только неподтвержденные сообщения.

Читать сообщения можно и [без подтверждения обработки](#no-commit). В этом случае при новом подключении будут прочитаны все неподтвержденные сообщения, в том числе и уже обработанные.

Информацию о том, какие сообщения уже обработаны, можно [сохранять на клиентской стороне](#client-commit), передавая на сервер стартовую позицию чтения при создании подключения. При этом позиция чтения сообщений на сервере не изменяется.

Можно использовать [транзакции](#read-tx). В этом случае позиция чтения изменится при подтверждении транзакции. При новом подключении будут прочитаны все неподтверждённые сообщения.

{% list tabs group=lang %}

- C++

  Работа пользователя с объектом `IReadSession` в общем устроена как обработка цикла событий со следующими типами событий: `TDataReceivedEvent`, `TCommitOffsetAcknowledgementEvent`, `TStartPartitionSessionEvent`, `TEndPartitionSessionEvent`, `TStopPartitionSessionEvent`, `TPartitionSessionStatusEvent`, `TPartitionSessionClosedEvent` и `TSessionClosedEvent`.

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

  Если обработчик для некоторого события не установлен, его необходимо получить и обработать в методах `GetEvent` / `GetEvents`. Для неблокирующего ожидания очередного события есть метод `WaitEvent` с сигнатурой `TFuture<void>()`.

- Go

  <!-- source: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->
  SDK получает данные с сервера партиями и буферизирует их. В зависимости от задач клиентский код может читать сообщения из буфера по одному или пакетами.
  <!-- endsource: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->

- Python

  <!-- source: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->
  SDK получает данные с сервера партиями и буферизирует их. В зависимости от задач клиентский код может читать сообщения из буфера по одному или пакетами.
  <!-- endsource: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->

- Java

  <!-- source: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->
  SDK получает данные с сервера партиями и буферизирует их. В зависимости от задач клиентский код может читать сообщения из буфера по одному или пакетами.
  <!-- endsource: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->

- C#

  <!-- source: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->
  SDK получает данные с сервера партиями и буферизирует их. В зависимости от задач клиентский код может читать сообщения из буфера по одному или пакетами.
  <!-- endsource: ru/reference/ydb-sdk/_includes/reading_messages_common.md -->

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  Полный пример чтения топика в транзакции с записью в таблицу: [`topic-read-in-transaction-example.rs`](https://github.com/ydb-platform/ydb-rs-sdk/blob/master/ydb/examples/topic-read-in-transaction-example.rs).

  ```rust
  let batch = reader.pop_batch_in_tx(&mut tx).await?;
  // обработка batch.messages и commit транзакции
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Чтение без подтверждения обработки сообщений {#no-commit}

#### Чтение сообщений по одному

{% list tabs group=lang %}

- C++

  Чтение сообщений по одному в C++ SDK не предусмотрено. Событие `TDataReceivedEvent` содержит пакет сообщений.

- Go

  ```go
  func SimpleReadMessages(ctx context.Context, r *topicreader.Reader) error {
      for {
          mess, err := r.ReadMessage(ctx)
          if err != nil {
              return err
          }
          processMessage(mess)
      }
  }
  ```

- Python

  {% list tabs %}

  - Native SDK

    ```python
    while True:
        message = reader.receive_message()
        process(message)
    ```

  - Native SDK (Asyncio)

    ```python
    while True:
        message = await reader.receive_message()
        process(message)
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    Чтобы читать сообщения без подтверждения обработки, по одному, используйте следующий код:

    ```java
    while(true) {
      Message message = reader.receive();
      process(message);
    }
    ```

  - Асинхронный API

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

  {% endlist %}

- C#

  ```c#
  try
  {
      while (!readerCts.IsCancellationRequested)
      {
          var message = await reader.ReadAsync(readerCts.Token);

          logger.LogInformation("Received message: [{MessageData}]", message.Data);
      }
  }
  catch (OperationCanceledException)
  {
  }
  ```

- JavaScript

  ```javascript
  for await (let batch of reader.read()) {
    for await (let msg of batch) {
    }
  }
  ```

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

#### Чтение сообщений пакетом

{% list tabs group=lang %}

- C++

  При установке сессии чтения с настройкой `SimpleDataHandlers` достаточно передать обработчик для сообщений с данными. SDK будет вызывать этот обработчик на каждый принятый от сервера пакет сообщений. Подтверждения чтения по умолчанию отправляться не будут.

  ```cpp
  auto settings = NYdb::NTopic::TReadSessionSettings()
      .EventHandlers_.SimpleDataHandlers(
          [](NYdb::NTopic::TReadSessionEvent::TDataReceivedEvent& event) {
              std::cout << "Get data event " << NYdb::NTopic::DebugString(event);
          }
      );

  auto session = topicClient.CreateReadSession(settings);

  // Wait SessionClosed event.
  session->GetEvent(/* block = */true);
  ```

  В этом примере после создания сессии основной поток дожидается завершения сессии со стороны сервера в методе `GetEvent`, другие типы событий приходить не будут.

- Go

  ```go
  func SimpleReadBatches(ctx context.Context, r *topicreader.Reader) error {
      for {
          batch, err := r.ReadMessageBatch(ctx)
          if err != nil {
              return err
          }
          processBatch(batch)
      }
  }
  ```

- Python

  {% list tabs %}

  - Native SDK

    ```python
    while True:
        batch = reader.receive_batch()
        process(batch)
    ```

  - Native SDK (Asyncio)

    ```python
    while True:
        batch = await reader.receive_batch()
        process(batch)
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    В синхронном клиенте нет возможности прочитать сразу пакет сообщений.

  - Асинхронный API

    Чтобы прочитать пакет сообщений без подтверждения обработки, используйте следующий код:

    ```java
    private class Handler extends AbstractReadEventHandler {
      @Override
      public void onMessages(DataReceivedEvent event) {
          for (Message message : event.getMessages()) {
              process(message);
          }
      }
    }
    ```

  {% endlist %}

- C#

  ```c#
  try
  {
      while (!readerCts.IsCancellationRequested)
      {
          var batchMessages = await reader.ReadBatchAsync(readerCts.Token);

          foreach (var message in batchMessages.Batch)
          {
              logger.LogInformation("Received message: [{MessageData}]", message.Data);
          }
      }
  }
  catch (OperationCanceledException)
  {
  }
  ```

- JavaScript

  ```javascript
  for await (let batch of reader.read()) {
  }
  ```

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Чтение с подтверждением обработки сообщений {#commit}

Подтверждение обработки сообщения (коммит) - сообщает серверу, что сообщение из топика обработано получателем и больше его отправлять не нужно. При использовании чтения с подтверждением нужно подтверждать все полученные сообщения без пропуска. Коммит сообщений на сервере происходит после подтверждения очередного интервала сообщений «без дырок», сами подтверждения при этом можно отправлять в любом порядке.

Например, с сервера пришли сообщения 1, 2, 3. Программа обрабатывает их параллельно и отправляет подтверждения в таком порядке: 1, 3, 2. В этом случае сначала будет закоммичено сообщение 1, а сообщения 2 и 3 будут закоммичены только после того как сервер получит подтверждение об обработке сообщения 2.

В случае ошибки на коммите сообщения можно написать эту ошибку в лог и продолжить работу. Состояние сообщения в этой точке неизвестно. Сообщение могло закоммититься, а потом возникла сетевая ошибка и клиент не получил подтверждения. Если сообщение не закоммитилось, то оно будет прочитано ещё раз и снова поступит в обработку (может быть на другом читателе). Ретраить именно коммит смысла нет, т.к. сессия чтения этого сообщения уже потеряна.

#### Чтение сообщений по одному с подтверждением

{% list tabs group=lang %}

- C++

  Чтение сообщений по одному в C++ SDK не предусмотрено. Событие `TDataReceivedEvent` содержит пакет сообщений.

- Go

  ```go
  func SimpleReadMessages(ctx context.Context, r *topicreader.Reader) error {
      for {
        mess, err := r.ReadMessage(ctx)
        if err != nil {
            return err
        }
        processMessage(mess)
        r.Commit(mess.Context(), mess)
      }
  }
  ```

  По умолчанию `Commit` — это быстрый вызов: сохраняет данные во внутреннем буфере и сразу возвращает управление, а реальная отправка происходит позже. Поэтому, чтобы не терять последние коммиты перед выходом из программы, читателя нужно закрывать явно с помощью вызова `Reader.Close()`.

- Python

  `commit` - это быстрый вызов: сохраняет данные во внутреннем буфере и сразу возвращает управление, а реальная отправка происходит позже. Поэтому, чтобы не терять последние коммиты перед выходом из программы, читателя нужно закрывать явно.

  {% list tabs %}

  - Native SDK

    ```python
    while True:
        message = reader.receive_message()
        process(message)
        reader.commit(message)
    ```

  - Native SDK (Asyncio)

    ```python
    while True:
        message = await reader.receive_message()
        process(message)
        reader.commit(message)
    ```

  {% endlist %}

- Java

  Для подтверждения обработки сообщения достаточно вызвать у сообщения метод `commit`.
  Актуально как для синхронного, так и для асинхронного читателя.
  В асинхронном читателе, при обработке пакета сообщений, можно вызвать `commit` или у всего пакета сразу, или у каждого сообщения отдельно.
  Этот метод возвращает `CompletableFuture<Void>`, успешное выполнение которой означает подтверждение обработки сервером.
  В случае ошибки коммита не следует пытаться его ретраить. Скорее всего, ошибка вызвана закрытием сессии.
  Читатель (необязательно этот же) сам создаст новую сессию для этой партиции и сообщение будет прочитано снова.

  ```java
  message.commit()
         .whenComplete((result, ex) -> {
             if (ex != null) {
                 // Read session was probably closed, there is nothing we can do here.
                 // Do not retry this commit on the same event.
                 logger.error("exception while committing message: ", ex);
             } else {
                 logger.info("message committed successfully");
             }
         });
  ```

- C#

  ```c#
  try
  {
      while (!readerCts.IsCancellationRequested)
      {
          var message = await reader.ReadAsync(readerCts.Token);

          logger.LogInformation("Received message: [{MessageData}]", message.Data);

          try
          {
              await message.CommitAsync();
          }
          catch (ReaderException e)
          {
              logger.LogError(e, "Failed to commit a message");
          }
      }
  }
  catch (OperationCanceledException)
  {
  }
  ```

- JavaScript

  ```javascript
  for await (let batch of reader.read()) {
    for (let msg of batch) {
      await reader.commit(msg);
    }
  }
  ```

- Rust

  ```rust
  let batch = reader.read_batch().await?;
  reader.commit(batch.get_commit_marker())?;
  // или с ожиданием ack от сервера:
  reader.commit_with_ack(batch.get_commit_marker()).await?;
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

#### Чтение сообщений пакетом с подтверждением

{% list tabs group=lang %}

- C++

  Аналогично [примеру выше](#no-commit), при установке сессии чтения с настройкой `SimpleDataHandlers` достаточно передать обработчик для сообщений с данными. SDK будет вызывать этот обработчик на каждый принятый от сервера пакет сообщений. Передача параметра `commitDataAfterProcessing = true` означает, что SDK будет отправлять на сервер подтверждения чтения всех сообщений после выполнения обработчика.

  ```cpp
  auto settings = NYdb::NTopic::TReadSessionSettings()
      .EventHandlers_.SimpleDataHandlers(
          [](NYdb::NTopic::TReadSessionEvent::TDataReceivedEvent& event) {
              std::cout << "Get data event " << NYdb::NTopic::DebugString(event);
          }
          , /* commitDataAfterProcessing = */true
      );

  auto session = topicClient.CreateReadSession(settings);

  // Wait SessionClosed event.
  session->GetEvent(/* block = */true);
  ```

- Go

  ```go
  func SimpleReadMessageBatch(ctx context.Context, r *topicreader.Reader) error {
      for {
        batch, err := r.ReadMessageBatch(ctx)
        if err != nil {
            return err
        }
        processBatch(batch)
        r.Commit(batch.Context(), batch)
      }
  }
  ```

  По умолчанию `Commit` - это быстрый вызов: сохраняет данные во внутреннем буфере и сразу возвращает управление, а реальная отправка происходит позже. Поэтому, чтобы не терять последние коммиты перед выходом из программы, читателя нужно закрывать явно.

- Python

  {% list tabs %}

  - Native SDK

    ```python
    while True:
        batch = reader.receive_batch()
        process(batch)
        reader.commit(batch)
    ```

  - Native SDK (Asyncio)

    ```python
    while True:
        batch = await reader.receive_batch()
        process(batch)
        reader.commit(batch)
    ```

  {% endlist %}

  `commit` - это быстрый вызов: сохраняет данные во внутреннем буфере и сразу возвращает управление, а реальная отправка происходит позже. Поэтому, чтобы не терять последние коммиты перед выходом из программы, читателя нужно закрывать явно.

- Java

  {% list tabs %}

  - Синхронный API

    Неактуально, т.к. в синхронном читателе нет возможности читать сообщения пакетами.

  - Асинхронный API

    В обработчике `onMessages` можно закоммитить весь пакет сообщений, вызвав `commit` на событии.

    ```java
    @Override
    public void onMessages(DataReceivedEvent event) {
      for (Message message : event.getMessages()) {
          process(message);
      }
      event.commit()
             .whenComplete((result, ex) -> {
                 if (ex != null) {
                     // Read session was probably closed, there is nothing we can do here.
                     // Do not retry this commit on the same message.
                     logger.error("exception while committing message batch: ", ex);
                 } else {
                     logger.info("message batch committed successfully");
                 }
             });
    }
    ```

  {% endlist %}

- C#

  ```c#
  try
  {
      while (!readerCts.IsCancellationRequested)
      {
          var batchMessages = await reader.ReadBatchAsync(readerCts.Token);

          foreach (var message in batchMessages.Batch)
          {
              logger.LogInformation("Received message: [{MessageData}]", message.Data);
          }

          try
          {
              await batchMessages.CommitBatchAsync();
          }
          catch (ReaderException e)
          {
              logger.LogError(e, "Failed to commit a message");
          }
      }
  }
  catch (OperationCanceledException)
  {
  }
  ```

- JavaScript

  ```javascript
  for await (let batch of reader.read()) {
    await reader.commit(batch);
  }
  ```

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Чтение с хранением позиции на клиентской стороне {#client-commit}

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

{% list tabs group=lang %}

- C++

  При обработке событий `TStartPartitionSessionEvent` можно при ответе серверу задать позицию, с которой следует начинать чтение.
  Для этого в метод `Confirm` следует передать параметр `readOffset`.
  Допольнительно можно передать параметр `commitOffset`, который укажет позицию, сообщения до которой следует считать [закоммиченными](#commit).

  Пример установки обработчика:

  ```cpp
  settings.EventHandlers_.StartPartitionSessionHandler(
      [](NYdb::NTopic::TReadSessionEvent::TStartPartitionSessionEvent& event) {
          auto readFromOffset = GetOffsetToReadFrom(event.GetPartitionId());
          event.Confirm(readFromOffset);
      }
  );
  ```

  Здесь `GetOffsetToReadFrom` - это часть примера, а не SDK. Используйте свой способ определить требуемую стартовую позицию чтения для партиции с данным partition id.

  Также в `TReadSessionSettings` поддерживается настройка `ReadFromTimestamp` для чтения событий с отметками времени записи не меньше данной. Эта настройка предполагается не для точного позиционирования старта, а для пропуска объёма данных за большой интервал времени. Несколько первых полученных сообщений могут иметь отметки времени записи меньше указанной.

- Go

  {% note tip %}

  В режиме читателя по умолчанию оффсеты до позиции, указанной через `res.StartFrom`, подтверждаются на сервере. После этого повторное чтение тех же сообщений путём сдвига позиции назад становится невозможным. Чтобы отключить автоматическое подтверждение, используйте режим без коммитов при создании читателя.

  ```go
  reader, err := db.Topic().StartReader(
    consumerName,
    topicoptions.ReadTopic(topicName),
    topicoptions.WithReaderCommitMode(topicoptions.CommitModeNone),
  )
  ```

  {% endnote %}

  ```go
  func ReadWithExplicitPartitionStartStopHandlerAndOwnReadProgressStorage(ctx context.Context, db ydb.Connection) error {
      readContext, stopReader := context.WithCancel(context.Background())
      defer stopReader()

      readStartPosition := func(
          ctx context.Context,
          req topicoptions.GetPartitionStartOffsetRequest,
      ) (res topicoptions.GetPartitionStartOffsetResponse, err error) {
          offset, err := readLastOffsetFromDB(ctx, req.Topic, req.PartitionID)
          res.StartFrom(offset)

          // Reader will stop if return err != nil
          return res, err
      }

      r, err := db.Topic().StartReader("my-consumer", topicoptions.ReadTopic("my-topic"),
          topicoptions.WithGetPartitionStartOffset(readStartPosition),
      )
      if err != nil {
          return err
      }

      go func() {
          <-readContext.Done()
          _ = r.Close(ctx)
      }()

      for {
          batch, err := r.ReadMessageBatch(readContext)
          if err != nil {
              return err
          }

          processBatch(batch)
          _ = externalSystemCommit(batch.Context(), batch.Topic(), batch.PartitionID(), batch.EndOffset())
      }
  }
  ```

- Python

  Функциональность находится в разработке.

- Java

  Чтение с заданного оффсета в Java возможно только в асинхронном читателе.
  В обработчике событий `StartPartitionSessionEvent` можно при ответе серверу задать позицию, с которой следует начинать чтение.
  Для этого в метод `confirm` следует передать настройки `StartPartitionSessionSettings` с указанным оффсетом через `setReadOffset`.
  Также вызовом `setCommitOffset` можно указать оффсет, который следует считать закоммиченным.

  ```java
  @Override
  public void onStartPartitionSession(StartPartitionSessionEvent event) {
      event.confirm(StartPartitionSessionSettings.newBuilder()
              .setReadOffset(lastReadOffset) // Long
              .setCommitOffset(lastCommitOffset) // Long
              .build());
  }
  ```

  Также поддерживается настройка читателя `setReadFrom` для чтения событий с отметками времени записи не меньше данной.

- JavaScript

  ```javascript
  await using reader = createTopicReader(driver, {
    topic: topicName,
    consumer: consumerName,
    onPartitionSessionStart: (evt) => {
      return {
        readOffset: 0n,
        commitOffset: 0n,
      };
    },
  });
  ```

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Чтение без указания Consumer'а {#no-consumer}

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

{% list tabs group=lang %}

- C++

  В `NYdb::NTopic::TReadSessionSettings` вызовите `WithoutConsumer()`:

  ```cpp
  auto settings = NYdb::NTopic::TReadSessionSettings()
      .WithoutConsumer()
      .AppendTopics(
          NYdb::NTopic::TTopicReadSettings("topic-path")
              .AppendPartitionIds(0)
              .AppendPartitionIds(1)
              .AppendPartitionIds(2));

  auto readSession = topicClient.CreateReadSession(settings);
  ```

  При переподключении прогресс чтения на сервере не сохраняется. Чтобы не начинать с начала, при каждом старте сессии чтения партиции передавайте смещение в `TStartPartitionSessionEvent::Confirm` — см. [хранение позиции на клиенте](#client-commit).

- Go

  Нужно передать пустую строку в качестве имени consumer и опцию `topicoptions.WithReaderWithoutConsumer(false)` (режим **экспериментальный**, см. [VERSIONING](https://github.com/ydb-platform/ydb-go-sdk/blob/master/VERSIONING.md) в репозитории SDK). В селекторе чтения укажите путь топика и список партиций. Коммиты сообщений в этом режиме недоступны (`CommitModeNone`); при переподключениях прогресс нужно восстанавливать на стороне клиента — см. [хранение позиции на клиенте](#client-commit).

  ```go
  reader, err := db.Topic().StartReader(
    "",
    topicoptions.ReadSelectors{{
      Path:       "topic-path",
      Partitions: []int64{0, 1, 2},
    }},
    topicoptions.WithReaderWithoutConsumer(false),
  )
  if err != nil {
    return err
  }
  ```

- Java

  Для чтения без Consumer'а следует в настройках читателя `ReaderSettings` это явно указать, вызвав `withoutConsumer()`:

  ```java
  ReaderSettings settings = ReaderSettings.newBuilder()
          .withoutConsumer()
          .addTopic(TopicReadSettings.newBuilder()
                  .setPath(TOPIC_NAME)
                  .build())
          .build();
  ```

  В таком случае нужно учитывать, что при переустановке соединения прогресс на сервере будет сброшен. Поэтому, чтобы не начинать чтение сначала, в SDK следует передавать offset начала чтения при каждом старте сессии чтения партиции:

  ```java
  @Override
  public void onStartPartitionSession(StartPartitionSessionEvent event) {
      event.confirm(StartPartitionSessionSettings.newBuilder()
              .setReadOffset(lastReadOffset) // the last offset read by this client, Long
              .build());
  }
  ```

- Python

  Для чтения без Consumer'а следует создать читателя с помощью метода `reader` с указанием следующих аргументов:
  - `topic` - объект `ydb.TopicReaderSelector` с указанными `path` и списком `partitions`;
  - `consumer` - должен быть `None`;
  - `event_handler` - наследник `ydb.TopicReaderEvents.EventHandler`, который реализует функцию `on_partition_get_start_offset`. Эта функция отвечает за возвращение начального смещения (offset) для чтения сообщений при старте читателя, а также во время переподключений. Клиентское приложение должно указать это смещение в параметре `ydb.TopicReaderEvents.OnPartitionGetStartOffsetResponse.start_offset`. Также функция может быть реализована как асинхронная.

  Пример:

  ```python
  class CustomEventHandler(ydb.TopicReaderEvents.EventHandler):
      def on_partition_get_start_offset(self, event: ydb.TopicReaderEvents.OnPartitionGetStartOffsetRequest):
          return ydb.TopicReaderEvents.OnPartitionGetStartOffsetResponse(
              start_offset=0,
          )

  reader = driver.topic_client.reader(
      topic=ydb.TopicReaderSelector(
          path="topic-path",
          partitions=[0, 1, 2],
      ),
      consumer=None,
      event_handler=CustomEventHandler(),
  )
  ```

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Чтение в транзакции {#read-tx}

{% list tabs group=lang %}

- C++

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

  [Пример на GitHub](https://github.com/ydb-platform/ydb-cpp-sdk/blob/main/examples/topic_reader/transaction/application.cpp)

  ```cpp
  readSession->WaitEvent().Wait(TDuration::Seconds(1));

  NYdb::NStatusHelpers::ThrowOnError(queryClient.RetryQuerySync([&readSession](NYdb::NQuery::TSession session) -> NYdb::TStatus {
      auto beginTxResult = session.BeginTransaction(NYdb::Query::TTxSettings::SerializableRW()).GetValueSync();
      if (!beginTxResult.IsSuccess()) {
          return beginTxResult;
      }
      auto tx = beginTxResult.GetTransaction();

      auto topicSettings = NYdb::NTopic::TReadSessionGetEventSettings()
          .Block(false);
          .Tx(tx);

      auto events = readSession->GetEvents(topicSettings);

      for (auto& event : events) {
          // обработать событие и записать результаты в таблицу
      }

      return tx.Commit().GetValueSync();
  }));
  ```

  {% note warning %}

  При обработке событий `events` не нужно явно подтверждать обработку для событий типа `TDataReceivedEvent`.

  {% endnote %}

  Подтверждение обработки события `TStopPartitionSessionEvent` надо делать после вызова `Commit`.

  ```cpp
  std::optional<NYdb::NTopic::TStopPartitionSessionEvent> stopPartitionSession;

  auto events = readSession->GetEvents(topicSettings);

  for (auto& event : events) {
      if (auto* e = std::get_if<NYdb::NTopic::TStopPartitionSessionEvent>(&event)) {
          stopPartitionSessionEvent = std::move(*e);
      } else {
          // обработать событие и записать результаты в таблицу
      }
  }

  auto commitResult = tx.Commit(commitSettings).GetValueSync();
  if (!commitResult.IsSuccess()) {
      return commitResult;
  }

  if (stopPartitionSessionEvent) {
      stopPartitionSessionEvent->Commit();
  }
  ```

- Go

  Для чтения сообщений в рамках транзакции следует использовать метод [`Reader.PopMessagesBatchTx`](https://pkg.go.dev/github.com/ydb-platform/ydb-go-sdk/v3/topic/topicreader#Reader.PopMessagesBatchTx). Он прочитает пакет сообщений и добавит их коммит в транзакцию, при этом отдельно коммитить эти сообщения не требуется. Читателя сообщений можно использовать повторно в разных транзакциях. При этом важно, чтобы порядок коммита транзакций соответствовал порядку получения сообщений от читателя, так как коммиты сообщений в топике должны выполняться строго по порядку. Проще всего это сделать если использовать читателя в цикле.

  [Пример на GitHub](https://github.com/ydb-platform/ydb-go-sdk/blob/master/examples/topic/topicreader/topic_reader_transaction.go)

  ```go
  for {
    err := db.Query().DoTx(ctx, func(ctx context.Context, tx query.TxActor) error {
      batch, err := reader.PopMessagesBatchTx(ctx, tx) // батч закоммитится при общем коммите транзакции
      if err != nil {
        return err
      }

      return processBatch(ctx, batch)
    })
    if err != nil {
      handleError(err)
    }
  }
  ```

- Python

  Для чтения сообщений в рамках транзакции следует использовать метод `reader.receive_batch_with_tx`. Он прочитает пакет сообщений и добавит их коммит в транзакцию, при этом отдельно коммитить эти сообщения не требуется. Читателя сообщений можно использовать повторно в разных транзакциях. При этом важно, чтобы порядок коммита транзакций соответствовал порядку получения сообщений от читателя, так как коммиты сообщений в топике должны выполняться строго по порядку - в противном случае транзакция получит ошибку на попытке сделать коммит. Проще всего это сделать, если использовать читателя в цикле.

  {% list tabs %}

  - Native SDK

    [Пример на GitHub](https://github.com/ydb-platform/ydb-python-sdk/blob/main/examples/topic/topic_transactions_example.py)

    ```python
    with driver.topic_client.reader(topic, consumer) as reader:
        with ydb.QuerySessionPool(driver) as session_pool:
            for _ in range(message_count):

                def callee(tx: ydb.QueryTxContext):
                    batch = reader.receive_batch_with_tx(tx, max_messages=1)
                    print(f"Message {batch.messages[0].data.decode()} was read with tx.")

                session_pool.retry_tx_sync(callee)
    ```

  - Native SDK (Asyncio)

    [Пример на GitHub](https://github.com/ydb-platform/ydb-python-sdk/blob/main/examples/topic/topic_transactions_async_example.py)

    ```python
    async with driver.topic_client.reader(topic, consumer) as reader:
        async with ydb.aio.QuerySessionPool(driver) as session_pool:
            for _ in range(message_count):

                async def callee(tx: ydb.aio.QueryTxContext):
                    batch = await reader.receive_batch_with_tx(tx, max_messages=1)
                    print(f"Message {batch.messages[0].data.decode()} was read with tx.")

                await session_pool.retry_tx_async(callee)
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    [Пример на GitHub](https://github.com/ydb-platform/ydb-java-examples/blob/develop/ydb-cookbook/src/main/java/tech/ydb/examples/topic/transactions/TransactionReadSync.java)

    В настройках `ReceiveSettings` метода `receive` можно указать транзакцию:

    ```java
    Message message = reader.receive(ReceiveSettings.newBuilder()
          .setTransaction(transaction)
          .build());
    ```

    Тогда полученное сообщение будет закоммичено вместе с транзакцией. Коммитить его отдельно не нужно.
    Метод `receive` свяжет на сервере оффсеты сообщения с транзакцией вызовом `sendUpdateOffsetsInTransaction` и вернёт управление, когда получит ответ на него.

  - Асинхронный API

    [Пример на GitHub](https://github.com/ydb-platform/ydb-java-examples/blob/develop/ydb-cookbook/src/main/java/tech/ydb/examples/topic/transactions/TransactionReadAsync.java)

    После получения сообщения в обработчике `onMessages` можно связать одно или несколько сообщений с транзакцией.
    Для этого нужно вызвать отдельный метод `reader.updateOffsetsInTransaction` и дождаться его выполнения на сервере.
    Этот метод принимает параметром список оффсетов. Для удобства у `Message` и `DataReceivedEvent` есть метод `getPartitionOffsets()`, возвращающий такой список.

    ```java
    @Override
    public void onMessages(DataReceivedEvent event) {
      for (Message message : event.getMessages()) {
          // creating a session in the table service
          Result<Session> sessionResult = tableClient.createSession(Duration.ofSeconds(10)).join();
          if (!sessionResult.isSuccess()) {
              logger.error("Couldn't get a session from the pool: {}", sessionResult);
              return; // retry or shutdown
          }
          Session session = sessionResult.getValue();
          // creating a transaction in the table service
          // this transaction is not yet active and has no id
          TableTransaction transaction = session.createNewTransaction(TxMode.SERIALIZABLE_RW);

          // do something else in the transaction
          transaction.executeDataQuery("SELECT 1").join();
          // now the transaction is active and has an id
          // analyzeQueryResultIfNeeded();

          Status updateStatus = reader.updateOffsetsInTransaction(transaction,
                          message.getPartitionOffsets(), new UpdateOffsetsInTransactionSettings.Builder().build())
                  // Do not commit a transaction without waiting for updateOffsetsInTransaction result to avoid a race condition
                  .join();
          if (!updateStatus.isSuccess()) {
              logger.error("Couldn't update offsets in a transaction: {}", updateStatus);
              return; // retry or shutdown
          }

          Status commitStatus = transaction.commit().join();
          analyzeCommitStatus(commitStatus);
      }
    }
    ```

  {% endlist %}

  <!-- source: ru/reference/ydb-sdk/_includes/alerts/java_transaction_requirements.md -->
  {% note info %}

  Требования к транзакции:

  * Это должна быть активная (имеющая идентификатор) транзакция в одном из сервисов YDB. Например, [Table](https://github.com/ydb-platform/ydb-java-sdk/blob/master/table/src/main/java/tech/ydb/table/transaction/TableTransaction.java) или [Query](https://github.com/ydb-platform/ydb-java-sdk/blob/master/query/src/main/java/tech/ydb/query/QueryTransaction.java).
  * В Topic Service поддерживается только уровень изоляции транзакций `SERIALIZABLE_RW`.

  {% endnote %}
  <!-- endsource: ru/reference/ydb-sdk/_includes/alerts/java_transaction_requirements.md -->

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  Полный пример чтения топика в транзакции с записью в таблицу: [`topic-read-in-transaction-example.rs`](https://github.com/ydb-platform/ydb-rs-sdk/blob/master/ydb/examples/topic-read-in-transaction-example.rs).

  ```rust
  let batch = reader.pop_batch_in_tx(&mut tx).await?;
  // обработка batch.messages и commit транзакции
  ```

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Обработка серверного прерывания чтения {#stop}

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

При *мягком прерывании* клиент получает уведомление, что сервер уже закончил отправку сообщений из партиции и больше сообщения читаться не будут. Клиент может завершить обработку сообщений и отправить подтверждение на сервер.

В случае *жесткого прерывания* клиент получает уведомление, что работать с сообщениями партиции больше нельзя. Клиент должен прекратить обработку прочитанных сообщений. Неподтвержденные сообщения будут переданы другому читателю.

#### Мягкое прерывание чтения {#soft-stop}

{% list tabs group=lang %}

- C++

  Мягкое прерывание приходит в виде события `TStopPartitionSessionEvent` с методом `Confirm`. Клиент может завершить обработку сообщений и отправить подтверждение на сервер.

  Фрагмент цикла событий может выглядеть так:

  ```cpp
  auto event = readSession->GetEvent(/*block=*/true);
  if (auto* stopPartitionSessionEvent = std::get_if<NYdb::NTopic::TReadSessionEvent::TStopPartitionSessionEvent>(&*event)) {
      stopPartitionSessionEvent->Confirm();
  } else {
    // other event types
  }
  ```

- Go

  Клиентский код сразу получает все имеющиеся в буфере (на стороне SDK) сообщения, даже если их не достаточно для формирования пакета при групповой обработке.

  ```go
  r, _ := db.Topic().StartReader("my-consumer", nil,
      topicoptions.WithBatchReadMinCount(1000),
  )

  for {
      batch, _ := r.ReadMessageBatch(ctx) // <- if partition soft stop batch can be less, then 1000
      processBatch(batch)
      _ = r.Commit(batch.Context(), batch)
  }
  ```

- Python

  Специальной обработки не требуется.

  {% list tabs %}

  - Native SDK

    ```python
    while True:
        batch = reader.receive_batch()
        process(batch)
        reader.commit(batch)
    ```

  - Native SDK (Asyncio)

    ```python
    while True:
        batch = await reader.receive_batch()
        process(batch)
        reader.commit(batch)
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    Неактуально, т.к. в синхронном читателе нет возможности настраивать обработку подобных событий.
    Клиент сразу ответит серверу подтверждением остановки.

  - Асинхронный API

    Для возможности реагировать на такое событие следует переопределить метод `onStopPartitionSession(StopPartitionSessionEvent event)` в объекте-наследнике `ReadEventHandler` (см [Подключение к топику для чтения сообщений](#start-reader)).
    `event.confirm()` обязательно должен быть вызван, т.к. сервер ожидает этого ответа для продолжения остановки.

    ```java
    @Override
    public void onStopPartitionSession(StopPartitionSessionEvent event) {
      logger.info("Partition session {} stopped. Committed offset: {}", event.getPartitionSessionId(),
              event.getCommittedOffset());
      // This event means that no more messages will be received by server
      // Received messages still can be read from ReaderBuffer
      // Messages still can be committed, until confirm() method is called

      // Confirm that session can be closed
      event.confirm();
    }
    ```

  {% endlist %}

- C#

  Специальной обработки не требуется.

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  Rust SDK обрабатывает события остановки и закрытия partition session внутренне; публичного API для настройки мягкого или жёсткого прерывания чтения пока нет.

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

#### Жесткое прерывание чтения {#hard-stop}

{% list tabs group=lang %}

- C++

  Жёсткое прерывание приходит в виде события `TPartitionSessionClosedEvent` либо в ответ на подтверждение мягкого прерывания, либо при потере соединения с партицией. Узнать причину можно, вызвав метод `GetReason`.

  Фрагмент цикла событий может выглядеть так:

  ```cpp
  auto event = readSession->GetEvent(/*block=*/true);
  if (auto* partitionSessionClosedEvent = std::get_if<NYdb::NTopic::TReadSessionEvent::TPartitionSessionClosedEvent>(&*event)) {
      if (partitionSessionClosedEvent->GetReason() == NYdb::NTopic::TPartitionSessionClosedEvent::EReason::ConnectionLost) {
          std::cout << "Connection with partition was lost" << std::endl;
      }
  } else {
    // other event types
  }
  ```

- Go

  При прерывании чтения контекст сообщения или пакета сообщений будет отменен.

  ```go
  ctx := batch.Context() // batch.Context() will cancel if partition revoke by server or connection broke
  if len(batch.Messages) == 0 {
      return
  }

  buf := &bytes.Buffer{}
  for _, mess := range batch.Messages {
      buf.Reset()
      _, _ = buf.ReadFrom(mess)
      _, _ = io.Copy(buf, mess)
      writeMessagesToDB(ctx, buf.Bytes())
  }
  ```

- Python

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

  {% list tabs %}

  - Native SDK

    ```python
    def process_batch(batch):
        for message in batch.messages:
            if not batch.alive:
                return False
            process(message)
        return True

    batch = reader.receive_batch()
    if process_batch(batch):
        reader.commit(batch)
    ```

  - Native SDK (Asyncio)

    ```python
    def process_batch(batch):
        for message in batch.messages:
            if not batch.alive:
                return False
            process(message)
        return True

    batch = await reader.receive_batch()
    if process_batch(batch):
        reader.commit(batch)
    ```

  {% endlist %}

- Java

  {% list tabs %}

  - Синхронный API

    Неактуально, т.к. в синхронном читателе нет возможности настраивать обработку подобных событий.

  - Асинхронный API

    ```java
    @Override
    public void onPartitionSessionClosed(PartitionSessionClosedEvent event) {
      logger.info("Partition session {} is closed.", event.getPartitionSession().getPartitionId());
    }
    ```

  {% endlist %}

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  Rust SDK обрабатывает события остановки и закрытия partition session внутренне; публичного API для настройки мягкого или жёсткого прерывания чтения пока нет.

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Поддержка автомасштабирования топиков {#autoscaling}

{% list tabs group=lang %}

- C++

  SDK поддерживает два режима чтения топиков с включенным автомасштабированием: режим полной поддержки и режим совместимости. Режим чтения задаётся в параметрах создания сессии чтения. По умолчанию используется режим совместимости.

  ```cpp
  auto settings = NYdb::NTopic::TReadSessionSettings()
      .SetAutoscalingSupport(true); // full support is enabled

  // or

  auto settings = NYdb::NTopic::TReadSessionSettings()
      .SetAutoscalingSupport(false); // compatibility mode is enabled

  auto readSession = topicClient.CreateReadSession(settings);
  ```

  В режиме полной поддержки, когда все сообщения из партиции будут прочитаны, придёт событие `TEndPartitionSessionEvent`. После получения этого события в партиции больше не появится новых сообщений для чтения. Чтобы продолжить чтение из дочерних партиций, необходимо вызвать `Confirm()`, тем самым подтвердив, что приложение готово принимать сообщения из дочерних партиций. Если сообщения из всех партиций обрабатываются в одном потоке, то `Confirm()` можно вызвать сразу после получения `TEndPartitionSessionEvent`. Если обработка сообщений из разных партиций осуществляется в разных потоках, то следует завершить обработку сообщений, например, выполнить накопившийся батч, подтвердить их обработку (коммит) или сохранить позицию чтения в своей базе, и только после этого вызвать `Confirm()`.

  После получения `TEndPartitionSessionEvent` и обработки всех сообщений рекомендуется всегда сразу подтверждать их обработку (коммит). Это позволит сбалансировать чтение дочерних партиций между разными сессиями чтения, что приведёт к равномерному распределению нагрузки по всем читателям.

  Фрагмент цикла событий может выглядеть так:

  ```cpp
  auto settings = NYdb::NTopic::TReadSessionSettings()
      .SetAutoscalingSupport(true);

  auto readSession = topicClient.CreateReadSession(settings);

  auto event = readSession->GetEvent(/*block=*/true);
  if (auto* endPartitionSessionEvent = std::get_if<NYdb::NTopic::TReadSessionEvent::TEndPartitionSessionEvent>(&*event)) {
      endPartitionSessionEvent->Confirm();
  } else {
    // other event types
  }
  ```

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

  Если клиент подтверждает обработку сообщений (коммит), то сигналом завершения обработки сообщений из партиции будет подтверждение обработки последнего сообщения этой партиции. В случае, если клиент не подтверждает обработку сообщений, сервер будет периодически прерывать чтение из партиции и переключаться на чтение в другой сессии (если существуют другие сессии, готовые обрабатывать партицию). Это будет продолжаться до тех пор, пока чтение не [начнётся](#client-commit) с конца партиции.

  Рекомендуется проверять корректность обработки мягкого прерывания чтения: клиент должен обработать полученные сообщения, подтвердить их обработку (коммит) или сохранить позицию чтения в своей базе, и только после этого вызывать `Confirm()` для события `TStopPartitionSessionEvent`.

- Go

  Включение автомасштабирования топика во время его создания производится с помощью опции `topicoptions.CreateWithAutoPartitioningSettings`:

  ```go
  import (
    ...

    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topicoptions"
    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topictypes"
  )

  err := db.Topic().Create(ctx,
    "topic",
    topicoptions.CreateWithAutoPartitioningSettings(
      topictypes.AutoPartitioningSettings{
        AutoPartitioningStrategy: topictypes.AutoPartitioningStrategyScaleUp,
      },
    ),
  )
  ```

  При необходимости в AutoPartitioningSettings можно задать и другие параметры:

  ```go
  err := db.Topic().Create(ctx,
    "topic",
    topicoptions.CreateWithAutoPartitioningSettings(
      topictypes.AutoPartitioningSettings{
        AutoPartitioningStrategy: topictypes.AutoPartitioningStrategyScaleUp,
        AutoPartitioningWriteSpeedStrategy: topictypes.AutoPartitioningWriteSpeedStrategy{
          StabilizationWindow:    time.Minute,
          UpUtilizationPercent:   80,
        },
      },
    ),
  )
  ```

  Включение автомасштабирования у существующего топика производится с помощью опции `topicoptions.AlterWithAutoPartitioningStrategy` у `.Topic().Alter`:

  ```go
  import (
    ...

    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topicoptions"
    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topictypes"
  )

  err := db.Topic().Alter(
    ctx,
    "topic",
    topicoptions.AlterWithAutoPartitioningStrategy(
      topictypes.AutoPartitioningStrategyScaleUp,
    ),
  )

  // другие опции
  err := db.Topic().Alter(
    ctx,
    "topic",
    topicoptions.AlterWithAutoPartitioningStrategy(
      topictypes.AutoPartitioningStrategyScaleUp,
    ),
    topicoptions.AlterWithAutoPartitioningWriteSpeedStabilizationWindow(time.Minute),
    topicoptions.AlterWithAutoPartitioningWriteSpeedUpUtilizationPercent(80),
  )
  ```

  SDK поддерживает два режима чтения топиков с включенным автомасштабированием: режим полной поддержки и режим совместимости. Режим чтения задаётся опцией `topicoptions.WithReaderSupportSplitMergePartitions` во время создания читателя. По умолчанию используется режим полной поддержки (`true`).

  ```go
  import (
    ...

    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topicoptions"
    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topictypes"
  )

  // режим полной поддержки (обработка автомасштабирования в SDK, по умолчанию)
  reader, err := db.Topic().StartReader(
    "consumer",
    topicoptions.ReadTopic("topic"),
    topicoptions.WithReaderSupportSplitMergePartitions(true),
  )

  // режим совместимости (обработка автомасштабирования на сервере)
  reader, err := db.Topic().StartReader(
    "consumer",
    topicoptions.ReadTopic("topic"),
    topicoptions.WithReaderSupportSplitMergePartitions(false),
  )
  ```

- Python

  Включение автомасштабирования топика во время его создания производится с помощью аргумента `auto_partitioning_settings` у `create_topic`:

  {% list tabs %}
  - Native SDK

    ```python
    driver.topic_client.create_topic(
        topic,
        consumers=[consumer],
        min_active_partitions=10,
        max_active_partitions=100,
        auto_partitioning_settings=ydb.TopicAutoPartitioningSettings(
            strategy=ydb.TopicAutoPartitioningStrategy.SCALE_UP,
            up_utilization_percent=80,
            down_utilization_percent=20,
            stabilization_window=datetime.timedelta(seconds=300),
        ),
    )
    ```

  - Native SDK (Asyncio)

    ```python
    await driver.topic_client.create_topic(
        topic,
        consumers=[consumer],
        min_active_partitions=10,
        max_active_partitions=100,
        auto_partitioning_settings=ydb.TopicAutoPartitioningSettings(
            strategy=ydb.TopicAutoPartitioningStrategy.SCALE_UP,
            up_utilization_percent=80,
            down_utilization_percent=20,
            stabilization_window=datetime.timedelta(seconds=300),
        ),
    )
    ```

  {% endlist %}

  Внесение изменений в существующий топик производится с помощью аргумента `alter_auto_partitioning_settings` у `alter_topic`:

  ```python
      driver.topic_client.alter_topic(
          topic_path,
          alter_auto_partitioning_settings=ydb.TopicAlterAutoPartitioningSettings(
              set_strategy=ydb.TopicAutoPartitioningStrategy.SCALE_UP,
              set_up_utilization_percent=80,
              set_down_utilization_percent=20,
              set_stabilization_window=datetime.timedelta(seconds=300),
          ),
      )
  ```

  SDK поддерживает два режима чтения топиков с включенным автомасштабированием: режим полной поддержки и режим совместимости. Режим чтения задаётся аргументом `auto_partitioning_support` во время создания читателя. По умолчанию используется режим полной поддержки.

  ```python
  reader = driver.topic_client.reader(
      topic,
      consumer,
      auto_partitioning_support=True, # Full support is enabled
  )

  # or

  reader = driver.topic_client.reader(
      topic,
      consumer,
      auto_partitioning_support=False, # Compatibility mode is enabled
  )
  ```

  С практической точки зрения для конечного пользователя режимы не отличаются. Режим полной поддержки отличается от режима совместимости тем, кто гарантирует порядок чтения — клиент или сервер. Режим совместимости достигается серверной обработкой и, как правило, работает медленнее.

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Java

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#311](https://github.com/ydb-platform/ydb-rs-sdk/issues/311)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}

### Подтверждение обработки вне читателя {#commit-outside-the-reader}

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

{% list tabs group=lang %}

- C++

  Подтверждение обработки вне сессии чтения производится с помощью метода `NYdb::NTopic::TTopicClient::CommitOffset`:

  ```cpp
  #include <ydb-cpp-sdk/client/topic/client.h>

  NYdb::NTopic::TTopicClient topicClient(driver);

  NYdb::NStatusHelpers::ThrowOnError(topicClient.CommitOffset(
      topicPath,
      partitionId,
      consumerName,
      offset).GetValueSync());
  ```

  Если в момент подтверждения существует активная сессия чтения (например через `CreateReadSession`), рекомендуется передать её идентификатор с помощью опции `ReadSessionId` в `NYdb::NTopic::TCommitOffsetSettings`. Это позволяет серверу не прерывать текущую сессию чтения:

  ```cpp
  // Получение идентификатора сессии чтения
  std::string sessionId = readSession->GetSessionId();

  NYdb::NStatusHelpers::ThrowOnError(topicClient.CommitOffset(
      topicPath,
      partitionId,
      consumerName,
      offset,
      NYdb::NTopic::TCommitOffsetSettings()
          .ReadSessionId(sessionId)
  ).GetValueSync());
  ```

- Go

  Подтверждение обработки вне читателя производится с помощью метода `db.Topic().CommitOffset`:

  ```go
  // Базовый способ — подтверждение оффсета без активной сессии чтения
  err := db.Topic().CommitOffset(
    ctx,
    topicPath,
    partitionID,
    consumer,
    offset,
  )
  ```

  Если в момент подтверждения существует активная сессия чтения (через `StartReader` или `StartListener`), рекомендуется передать её идентификатор с помощью опции `WithCommitOffsetReadSessionID`. Это позволяет серверу не прерывать текущую сессию чтения:

  ```go
  import (
    // ...
    "github.com/ydb-platform/ydb-go-sdk/v3/topic/topicoptions"
  )

  // Получение идентификатора сессии чтения
  sessionID := reader.ReadSessionID()
  // или: sessionID := listener.ReadSessionID()

  err = db.Topic().CommitOffset(
    ctx,
    topicPath,
    partitionID,
    consumer,
    offset,
    topicoptions.WithCommitOffsetReadSessionID(sessionID),
  )
  ```

- Python

  Подтверждение обработки вне читателя производится с помощью метода `topic_client.commit_offset`:

  {% list tabs %}

  - Native SDK

    ```python
    driver.topic_client.commit_offset(
        topic_path,
        consumer_name,
        partition_id,
        offset,
        reader.read_session_id,  # опционально: не прерывает активную сессию чтения
    )
    ```

  - Native SDK (Asyncio)

    ```python
    await driver.topic_client.commit_offset(
        topic_path,
        consumer_name,
        partition_id,
        offset,
        reader.read_session_id,  # опционально: не прерывает активную сессию чтения
    )
    ```

  {% endlist %}

- JavaScript

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Java

  ```java
  TopicClient client = ...;

  String sessionID = reader.getSessionId();
  // У AsyncReader идентификатор сессии можно получить при обработке события SessionStartedEvent

  client.commitOffset(
      topicPath,
      CommitOffsetSettings.newBuilder()
          .setReadSessionId(sessionID)
          .setPartitionId(partitionID)
          .setConsumer(consumer)
          .setOffset(offset)
          .build()
  ).join().expectSuccess("Error commit!");
  ```

- C#

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

- Rust

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

  Отслеживать прогресс или проголосовать за поддержку в Rust SDK: [ydb-rs-sdk#330](https://github.com/ydb-platform/ydb-rs-sdk/issues/330)

- PHP

  <!-- source: ru/_includes/feature-not-supported.md -->
  Функциональность на данный момент не поддерживается.
  <!-- endsource: ru/_includes/feature-not-supported.md -->

{% endlist %}
