> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-parallel-read-in-order-multi-part.mintlify.site/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](/ru/concepts/features/operations/insert/deduplication) и [Deduplicate insertion config](/ru/reference/settings/session-settings/insert#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).
