> ## Documentation Index
> Fetch the complete documentation index at: https://clickhouse.com/docs/llms.txt
> Use this file to discover all available pages before exploring further.

> Вы можете загружать данные в ClickHouse с помощью Apache Beam

# Интеграция Apache Beam с ClickHouse

export const ClickHouseSupportedBadge = () => {
  return <div className="ClickHouseSupportedBadge">
            <div className="ClickHouseSupportedIcon">
                <svg width="16" height="16" viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
                    <path d="M1.30762 1.39073C1.30762 1.3103 1.37465 1.22986 1.46849 1.22986H2.64824C2.72868 1.22986 2.80912 1.29689 2.80912 1.39073V14.4886C2.80912 14.5691 2.74209 14.6495 2.64824 14.6495H1.46849C1.38805 14.6495 1.30762 14.5825 1.30762 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M4.2832 1.39073C4.2832 1.3103 4.35023 1.22986 4.44408 1.22986H5.62383C5.70427 1.22986 5.7847 1.29689 5.7847 1.39073V14.4886C5.7847 14.5691 5.71767 14.6495 5.62383 14.6495H4.44408C4.36364 14.6495 4.2832 14.5825 4.2832 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M7.25977 1.39073C7.25977 1.3103 7.3268 1.22986 7.42064 1.22986H8.60039C8.68083 1.22986 8.76127 1.29689 8.76127 1.39073V14.4886C8.76127 14.5691 8.69423 14.6495 8.60039 14.6495H7.42064C7.3402 14.6495 7.25977 14.5825 7.25977 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M10.2354 1.39073C10.2354 1.3103 10.3024 1.22986 10.3962 1.22986H11.576C11.6564 1.22986 11.7369 1.29689 11.7369 1.39073V14.4886C11.7369 14.5691 11.6698 14.6495 11.576 14.6495H10.3962C10.3158 14.6495 10.2354 14.5825 10.2354 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M13.2256 6.6057C13.2256 6.52526 13.2926 6.44482 13.3865 6.44482H14.5662C14.6466 6.44482 14.7271 6.51186 14.7271 6.6057V9.27354C14.7271 9.35398 14.6601 9.43442 14.5662 9.43442H13.3865C13.306 9.43442 13.2256 9.36739 13.2256 9.27354V6.6057Z" fill="currentColor" />
                </svg>
            </div>
            Поддерживается в ClickHouse
        </div>;
};

<ClickHouseSupportedBadge />

**Apache Beam** — это открытая унифицированная модель программирования, которая позволяет разработчикам определять и выполнять как батч-, так и потоковые (непрерывные) конвейеры обработки данных. Гибкость Apache Beam заключается в поддержке широкого спектра сценариев обработки данных — от операций ETL (извлечение, преобразование, загрузка) до сложной обработки событий и Real-time аналитики.
Эта интеграция использует официальный [коннектор JDBC](https://github.com/ClickHouse/clickhouse-java) ClickHouse в качестве базового слоя для вставки.

<div id="integration-package">
  ## Пакет интеграции
</div>

Пакет интеграции, необходимый для работы Apache Beam с ClickHouse, поддерживается и разрабатывается в рамках [Apache Beam I/O Connectors](https://beam.apache.org/documentation/io/connectors/) — набора интеграций для множества популярных систем хранения данных и баз данных.
Реализация `org.apache.beam.sdk.io.clickhouse.ClickHouseIO` находится в [репозитории Apache Beam](https://github.com/apache/beam/tree/0bf43078130d7a258a0f1638a921d6d5287ca01e/sdks/java/io/clickhouse/src/main/java/org/apache/beam/sdk/io/clickhouse).

<div id="setup-of-the-apache-beam-clickhouse-package">
  ## Настройка пакета ClickHouse для Apache Beam
</div>

<div id="package-installation">
  ### Установка пакета
</div>

Добавьте следующую зависимость в систему управления пакетами:

```xml theme={null}
<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-clickhouse</artifactId>
    <version>${beam.version}</version>
</dependency>
```

<Warning>
  **Рекомендуемая версия Beam**

  Коннектор `ClickHouseIO` рекомендуется использовать с Apache Beam версии `2.59.0` и выше.
  Более ранние версии могут не полностью поддерживать функциональность коннектора.
</Warning>

Артефакты доступны в [официальном репозитории Maven](https://mvnrepository.com/artifact/org.apache.beam/beam-sdks-java-io-clickhouse).

<div id="code-example">
  ### Пример кода
</div>

В следующем примере CSV-файл `input.csv` считывается в виде `PCollection`, преобразуется в объект `Row` (с использованием заданной схемы) и вставляется в локальный экземпляр ClickHouse с помощью `ClickHouseIO`:

```java theme={null}

package org.example;

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.io.clickhouse.ClickHouseIO;
import org.apache.beam.sdk.schemas.Schema;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;
import org.joda.time.DateTime;

public class Main {

    public static void main(String[] args) {
        // Создание объекта Pipeline.
        Pipeline p = Pipeline.create();

        Schema SCHEMA =
                Schema.builder()
                        .addField(Schema.Field.of("name", Schema.FieldType.STRING).withNullable(true))
                        .addField(Schema.Field.of("age", Schema.FieldType.INT16).withNullable(true))
                        .addField(Schema.Field.of("insertion_time", Schema.FieldType.DATETIME).withNullable(false))
                        .build();

        // Применение преобразований к конвейеру.
        PCollection<String> lines = p.apply("ReadLines", TextIO.read().from("src/main/resources/input.csv"));

        PCollection<Row> rows = lines.apply("ConvertToRow", ParDo.of(new DoFn<String, Row>() {
            @ProcessElement
            public void processElement(@Element String line, OutputReceiver<Row> out) {

                String[] values = line.split(",");
                Row row = Row.withSchema(SCHEMA)
                        .addValues(values[0], Short.parseShort(values[1]), DateTime.now())
                        .build();
                out.output(row);
            }
        })).setRowSchema(SCHEMA);

        rows.apply("Write to ClickHouse",
                        ClickHouseIO.write("jdbc:clickhouse://localhost:8123/default?user=default&password=******", "test_table"));

        // Запуск конвейера.
        p.run().waitUntilFinish();
    }
}

```

<div id="supported-data-types">
  ## Поддерживаемые типы данных
</div>

| ClickHouse                         | Apache Beam                | Поддерживается | Примечания                                                                                                                                                       |
| ---------------------------------- | -------------------------- | -------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `TableSchema.TypeName.FLOAT32`     | `Schema.TypeName#FLOAT`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.FLOAT64`     | `Schema.TypeName#DOUBLE`   | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.INT8`        | `Schema.TypeName#BYTE`     | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.INT16`       | `Schema.TypeName#INT16`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.INT32`       | `Schema.TypeName#INT32`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.INT64`       | `Schema.TypeName#INT64`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.STRING`      | `Schema.TypeName#STRING`   | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.UINT8`       | `Schema.TypeName#INT16`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.UINT16`      | `Schema.TypeName#INT32`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.UINT32`      | `Schema.TypeName#INT64`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.UINT64`      | `Schema.TypeName#INT64`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.DATE`        | `Schema.TypeName#DATETIME` | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.DATETIME`    | `Schema.TypeName#DATETIME` | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.ARRAY`       | `Schema.TypeName#ARRAY`    | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.ENUM8`       | `Schema.TypeName#STRING`   | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.ENUM16`      | `Schema.TypeName#STRING`   | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.BOOL`        | `Schema.TypeName#BOOLEAN`  | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.TUPLE`       | `Schema.TypeName#ROW`      | ✅              |                                                                                                                                                                  |
| `TableSchema.TypeName.FIXEDSTRING` | `FixedBytes`               | ✅              | `FixedBytes` — это `LogicalType`, представляющий массив <br /> байтов фиксированной длины, который находится в <br /> `org.apache.beam.sdk.schemas.logicaltypes` |
|                                    | `Schema.TypeName#DECIMAL`  | ❌              |                                                                                                                                                                  |
|                                    | `Schema.TypeName#MAP`      | ❌              |                                                                                                                                                                  |

<div id="clickhouseiowrite-parameters">
  ## Параметры ClickHouseIO.Write
</div>

Конфигурацию `ClickHouseIO.Write` можно настроить с помощью следующих функций-сеттеров:

| Функция-сеттер параметра    | Тип аргумента               | Значение по умолчанию         | Описание                                                             |
| --------------------------- | --------------------------- | ----------------------------- | -------------------------------------------------------------------- |
| `withMaxInsertBlockSize`    | `(long maxInsertBlockSize)` | `1000000`                     | Максимальный размер блока строк для вставки.                         |
| `withMaxRetries`            | `(int maxRetries)`          | `5`                           | Максимальное количество повторных попыток для неудачных вставок.     |
| `withMaxCumulativeBackoff`  | `(Duration maxBackoff)`     | `Duration.standardDays(1000)` | Максимальная суммарная длительность задержки для повторных попыток.  |
| `withInitialBackoff`        | `(Duration initialBackoff)` | `Duration.standardSeconds(5)` | Начальная длительность задержки перед первой повторной попыткой.     |
| `withInsertDistributedSync` | `(Boolean sync)`            | `true`                        | Если `true`, синхронизирует операции вставки для distributed таблиц. |
| `withInsertQuorum`          | `(Long quorum)`             | `null`                        | Количество реплик, необходимых для подтверждения операции вставки.   |
| `withInsertDeduplicate`     | `(Boolean deduplicate)`     | `true`                        | Если `true`, для операций вставки включается дедупликация.           |
| `withTableSchema`           | `(TableSchema schema)`      | `null`                        | Схема целевой таблицы ClickHouse.                                    |

<div id="limitations">
  ## Ограничения
</div>

Учитывайте следующие ограничения при использовании коннектора:

* На данный момент поддерживается только операция Sink. Коннектор не поддерживает операцию Source.
* ClickHouse выполняет дедупликацию при вставке в таблицу `ReplicatedMergeTree` или в таблицу `Distributed`, построенную поверх `ReplicatedMergeTree`. Без репликации вставка в обычную таблицу MergeTree может приводить к появлению дубликатов, если вставка завершается ошибкой, а затем успешно повторяется. Однако каждый блок вставляется атомарно, а размер блока можно настроить с помощью `ClickHouseIO.Write.withMaxInsertBlockSize(long)`. Дедупликация достигается за счёт использования контрольных сумм вставленных блоков. Подробнее о дедупликации см. в разделах [Deduplication](/docs/ru/concepts/features/operations/insert/deduplication) и [Deduplicate insertion config](/docs/ru/reference/settings/session-settings#insert_deduplicate).
* Коннектор не выполняет никаких DDL-операторов; поэтому целевая таблица должна существовать до вставки.

<div id="related-content">
  ## Связанные материалы
</div>

* Документация по классу `ClickHouseIO`: [documentation](https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/clickhouse/ClickHouseIO.html).
* Репозиторий GitHub с примерами: [clickhouse-beam-connector](https://github.com/ClickHouse/clickhouse-beam-connector).
