> ## 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.

> Использование движка таблицы Kafka

# Использование движка таблицы Kafka

export const Image = ({img, alt, size = "lg"}) => {
  const normalizedSize = ["sm", "md", "lg"].includes(size) ? size : "lg";
  return <div className={`ch-image-${normalizedSize}`}>
      <Frame>
        <img src={img} alt={alt} />
      </Frame>
    </div>;
};

Движок таблицы Kafka можно использовать для [**чтения** данных из](#kafka-to-clickhouse) и [**записи** данных в](#clickhouse-to-kafka) Apache Kafka и другие брокеры с поддержкой Kafka API (например, Redpanda, Amazon MSK).

<div id="kafka-to-clickhouse">
  ### Kafka в ClickHouse
</div>

<Note>
  Если вы используете ClickHouse Cloud, мы рекомендуем вместо этого [ClickPipes](/docs/ru/integrations/clickpipes/home). ClickPipes изначально поддерживает подключения к частным сетям, независимое масштабирование ресурсов ингестии и кластера, а также комплексный мониторинг при стриминге данных из Kafka в ClickHouse.
</Note>

Для использования движка таблицы Kafka вам следует в общих чертах понимать, как работают [materialized views ClickHouse](/docs/ru/concepts/features/materialized-views/cascading-materialized-views).

<div id="overview">
  #### Обзор
</div>

Сначала мы рассмотрим наиболее распространённый сценарий использования: применение движка таблицы Kafka для вставки данных из Kafka в ClickHouse.

Движок таблицы Kafka позволяет ClickHouse напрямую читать из топика Kafka. Хотя это удобно для просмотра сообщений в топике, по своей конструкции этот движок допускает только однократное чтение: когда к таблице выполняется запрос, он считывает данные из очереди и увеличивает смещение потребителя перед возвратом результатов вызывающей стороне. На практике данные нельзя прочитать повторно без сброса этих смещений.

Чтобы сохранить данные, прочитанные через движок таблицы, нужен способ захватить их и вставить в другую таблицу. Эту возможность нативно предоставляют materialized view на основе триггеров. Materialized view инициирует чтение из движка таблицы, получая батчи документов. Предложение TO определяет пункт назначения данных — обычно это таблица семейства [MergeTree](/docs/ru/reference/engines/table-engines/mergetree-family/index). Этот процесс показан ниже:

<Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/kafka_01.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=fd6990e133e6b46eb62f7057956fb5a3" size="lg" alt="Диаграмма архитектуры движка таблицы Kafka" style={{width: '80%'}} width="2048" height="852" data-path="images/integrations/data-ingestion/kafka/kafka_01.webp" />

<div id="steps">
  #### Шаги
</div>

<Steps>
  <Step title="Подготовьте" id="1-prepare">
    Если у вас есть данные, загруженные в целевой топик, вы можете адаптировать приведённый ниже пример для использования со своим набором данных. Либо можно воспользоваться примером набора данных GitHub, доступным [здесь](https://datasets-documentation.s3.eu-west-3.amazonaws.com/kafka/github_all_columns.ndjson). Этот набор данных используется в примерах ниже и для краткости включает сокращённую схему и подмножество строк (в частности, мы ограничиваемся событиями GitHub, относящимися к [репозиторию ClickHouse](https://github.com/ClickHouse/ClickHouse)) по сравнению с полным набором данных, доступным [здесь](https://ghe.clickhouse.tech/). Тем не менее его достаточно, чтобы работало большинство [запросов, опубликованных вместе с набором данных](https://ghe.clickhouse.tech/).
  </Step>

  <Step title="Настройте ClickHouse" id="2-configure-clickhouse">
    Этот шаг обязателен, если вы подключаетесь к защищённому кластеру Kafka. Эти настройки нельзя передать через команды SQL DDL, их необходимо задать в `config.xml` ClickHouse. Предполагается, что вы подключаетесь к экземпляру, защищённому с помощью SASL. Это самый простой способ при работе с Confluent Cloud.

    ```xml theme={null}
    <clickhouse>
       <kafka>
           <sasl_username>username</sasl_username>
           <sasl_password>password</sasl_password>
           <security_protocol>sasl_ssl</security_protocol>
           <sasl_mechanisms>PLAIN</sasl_mechanisms>
       </kafka>
    </clickhouse>
    ```

    Либо поместите приведённый выше фрагмент в новый файл в каталоге conf.d/, либо объедините его с существующими файлами конфигурации. Сведения о доступных для настройки параметрах см. [здесь](/docs/ru/reference/engines/table-engines/integrations/kafka#configuration).

    Мы также создадим базу данных с именем `KafkaEngine`, которую будем использовать в этом руководстве:

    ```sql theme={null}
    CREATE DATABASE KafkaEngine;
    ```

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

    ```sql theme={null}
    USE KafkaEngine;
    ```
  </Step>

  <Step title="Создайте целевую таблицу" id="3-create-the-destination-table">
    Подготовьте целевую таблицу. В примере ниже для краткости используется сокращённая схема GitHub. Обратите внимание: хотя здесь используется движок таблицы MergeTree, этот пример можно легко адаптировать для любого движка из [семейства MergeTree](/docs/ru/reference/engines/table-engines/mergetree-family/index).

    ```sql theme={null}
    CREATE TABLE github
    (
        file_time DateTime,
        event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4, 'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
        actor_login LowCardinality(String),
        repo_name LowCardinality(String),
        created_at DateTime,
        updated_at DateTime,
        action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
        comment_id UInt64,
        path String,
        ref LowCardinality(String),
        ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
        creator_user_login LowCardinality(String),
        number UInt32,
        title String,
        labels Array(LowCardinality(String)),
        state Enum('none' = 0, 'open' = 1, 'closed' = 2),
        assignee LowCardinality(String),
        assignees Array(LowCardinality(String)),
        closed_at DateTime,
        merged_at DateTime,
        merge_commit_sha String,
        requested_reviewers Array(LowCardinality(String)),
        merged_by LowCardinality(String),
        review_comments UInt32,
        member_login LowCardinality(String)
    ) ENGINE = MergeTree ORDER BY (event_type, repo_name, created_at)
    ```
  </Step>

  <Step title="Создайте топик и заполните его" id="4-create-and-populate-the-topic">
    Далее мы создадим топик. Для этого можно использовать несколько инструментов. Если Kafka запущен локально на вашей машине или в контейнере Docker, хорошо подойдёт [RPK](https://docs.redpanda.com/current/get-started/rpk-install/). Мы можем создать топик с именем `github` с 5 партициями, выполнив следующую команду:

    ```bash theme={null}
    rpk topic create -p 5 github --brokers <host>:<port>
    ```

    Если мы запускаем Kafka в Confluent Cloud, то можем предпочесть использовать [Confluent CLI](https://docs.confluent.io/platform/current/tutorials/examples/clients/docs/kcat.html#produce-records):

    ```bash theme={null}
    confluent kafka topic create --if-not-exists github
    ```

    Теперь нужно заполнить этот topic данными; для этого мы будем использовать [kcat](https://github.com/edenhill/kcat). Если Kafka запущен локально и аутентификация отключена, можно выполнить команду, аналогичную следующей:

    ```bash theme={null}
    cat github_all_columns.ndjson |
    kcat -P \
      -b <host>:<port> \
      -t github
    ```

    Или следующий вариант, если в нашем кластере Kafka для аутентификации используется SASL:

    ```bash theme={null}
    cat github_all_columns.ndjson |
    kcat -P \
      -b <host>:<port> \
      -t github
      -X security.protocol=sasl_ssl \
      -X sasl.mechanisms=PLAIN \
      -X sasl.username=<username>  \
      -X sasl.password=<password> \
    ```

    Датасет содержит 200 000 строк, поэтому его приём займёт всего несколько секунд. Если вы хотите работать с более крупным датасетом, ознакомьтесь с [разделом о больших датасетах](https://github.com/ClickHouse/kafka-samples/tree/main/producer#large-datasets) в GitHub-репозитории [ClickHouse/kafka-samples](https://github.com/ClickHouse/kafka-samples).
  </Step>

  <Step title="Создайте движок таблицы Kafka" id="5-create-the-kafka-table-engine">
    В примере ниже создаётся движок таблицы с той же схемой, что и у таблицы MergeTree. Это не является строгим требованием, так как в целевой таблице могут быть псевдонимы или эфемерные столбцы. Однако настройки важны — обратите внимание на использование `JSONEachRow` в качестве типа данных для чтения JSON из топика Kafka. Значения `github` и `clickhouse` представляют собой имена топика и группы потребителей соответственно. На самом деле топики могут быть списком значений.

    ```sql theme={null}
    CREATE TABLE github_queue
    (
        file_time DateTime,
        event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4, 'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
        actor_login LowCardinality(String),
        repo_name LowCardinality(String),
        created_at DateTime,
        updated_at DateTime,
        action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
        comment_id UInt64,
        path String,
        ref LowCardinality(String),
        ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
        creator_user_login LowCardinality(String),
        number UInt32,
        title String,
        labels Array(LowCardinality(String)),
        state Enum('none' = 0, 'open' = 1, 'closed' = 2),
        assignee LowCardinality(String),
        assignees Array(LowCardinality(String)),
        closed_at DateTime,
        merged_at DateTime,
        merge_commit_sha String,
        requested_reviewers Array(LowCardinality(String)),
        merged_by LowCardinality(String),
        review_comments UInt32,
        member_login LowCardinality(String)
    )
       ENGINE = Kafka('kafka_host:9092', 'github', 'clickhouse',
                'JSONEachRow') SETTINGS kafka_thread_per_consumer = 0, kafka_num_consumers = 1;
    ```

    Ниже мы рассмотрим настройки движка и настройку производительности. На этом этапе простой запрос select к таблице `github_queue` должен прочитать несколько строк.  Обратите внимание, что это сдвинет смещения потребителя вперёд, из-за чего эти строки нельзя будет прочитать повторно без [сброса](#common-operations). Обратите внимание на ограничение и обязательный параметр `stream_like_engine_allow_direct_select.`
  </Step>

  <Step title="Создайте materialized view" id="6-create-the-materialized-view">
    materialized view свяжет две ранее созданные таблицы, считывая данные из таблицы Kafka и вставляя их в целевую таблицу MergeTree. Мы можем выполнить ряд преобразований данных. Здесь мы ограничимся простым чтением и вставкой. Использование \* предполагает, что имена столбцов совпадают (с учётом регистра).

    ```sql theme={null}
    CREATE MATERIALIZED VIEW github_mv TO github AS
    SELECT *
    FROM github_queue;
    ```

    В момент создания materialized view подключается к движку Kafka и начинает читать данные, вставляя строки в целевую таблицу. Этот процесс будет продолжаться бесконечно, при этом новые сообщения, вставляемые в Kafka, будут потребляться. При необходимости вы можете повторно запустить скрипт вставки, чтобы добавить в Kafka дополнительные сообщения.
  </Step>

  <Step title="Убедитесь, что строки вставлены" id="7-confirm-rows-have-been-inserted">
    Убедитесь, что в целевой таблице есть данные:

    ```sql theme={null}
    SELECT count() FROM github;
    ```

    Вы должны увидеть 200 000 строк:

    ```response theme={null}
    ┌─count()─┐
    │  200000 │
    └─────────┘
    ```
  </Step>
</Steps>

<div id="common-operations">
  #### Типовые операции
</div>

<div id="stopping--restarting-message-consumption">
  ##### Остановка и возобновление потребления сообщений
</div>

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

```sql theme={null}
DETACH TABLE github_queue;
```

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

```sql theme={null}
ATTACH TABLE github_queue;
```

<div id="adding-kafka-metadata">
  ##### Добавление метаданных Kafka
</div>

Может быть полезно отслеживать метаданные исходных сообщений Kafka после их приёма в ClickHouse. Например, может понадобиться узнать, какую часть конкретного топика или партиции мы уже обработали. Для этого движок таблицы Kafka предоставляет несколько [виртуальных столбцов](/docs/ru/reference/engines/table-engines/index#table_engines-virtual_columns). Их можно сохранить в виде столбцов в целевой таблице, изменив схему и оператор SELECT в materialized view.

Сначала выполните описанную выше операцию остановки, а затем добавьте столбцы в целевую таблицу.

```sql theme={null}
DETACH TABLE github_queue;
```

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

```sql theme={null}
ALTER TABLE github
   ADD COLUMN topic String,
   ADD COLUMN partition UInt64;
```

Далее нужно убедиться, что виртуальные столбцы сопоставлены должным образом.
Виртуальные столбцы имеют префикс `_`.
Полный список виртуальных столбцов можно найти [здесь](/docs/ru/reference/engines/table-engines/integrations/kafka#virtual-columns).

Чтобы обновить таблицу, добавив виртуальные столбцы, потребуется удалить materialized view, повторно выполнить ATTACH для таблицы с движком Kafka и заново создать materialized view.

```sql theme={null}
DROP VIEW github_mv;
```

```sql theme={null}
ATTACH TABLE github_queue;
```

```sql theme={null}
CREATE MATERIALIZED VIEW github_mv TO github AS
SELECT *, _topic AS topic, _partition as partition
FROM github_queue;
```

У недавно прочитанных строк должны быть метаданные.

```sql theme={null}
SELECT actor_login, event_type, created_at, topic, partition
FROM github
LIMIT 10;
```

Результат будет выглядеть так:

| actor\_login  | event\_type        | created\_at         | topic  | partition |
| :------------ | :----------------- | :------------------ | :----- | :-------- |
| IgorMinar     | CommitCommentEvent | 2011-02-12 02:22:00 | github | 0         |
| queeup        | CommitCommentEvent | 2011-02-12 02:23:23 | github | 0         |
| IgorMinar     | CommitCommentEvent | 2011-02-12 02:23:24 | github | 0         |
| IgorMinar     | CommitCommentEvent | 2011-02-12 02:24:50 | github | 0         |
| IgorMinar     | CommitCommentEvent | 2011-02-12 02:25:20 | github | 0         |
| dapi          | CommitCommentEvent | 2011-02-12 06:18:36 | github | 0         |
| sourcerebels  | CommitCommentEvent | 2011-02-12 06:34:10 | github | 0         |
| jamierumbelow | CommitCommentEvent | 2011-02-12 12:21:40 | github | 0         |
| jpn           | CommitCommentEvent | 2011-02-12 12:24:31 | github | 0         |
| Oxonium       | CommitCommentEvent | 2011-02-12 12:31:28 | github | 0         |

<div id="modify-kafka-engine-settings">
  ##### Изменение настроек движка Kafka
</div>

Мы рекомендуем удалить таблицу с движком Kafka и заново создать её с новыми настройками. В ходе этого процесса materialized view изменять не нужно — потребление сообщений возобновится, как только таблица с движком Kafka будет создана заново.

<div id="debugging-issues">
  ##### Отладка проблем
</div>

Ошибки, например связанные с аутентификацией, не отображаются в ответах на DDL-запросы к движку Kafka. Для диагностики проблем рекомендуем использовать основной файл журнала ClickHouse `clickhouse-server.err.log`. Дополнительное журналирование трассировки для базовой клиентской библиотеки Kafka [librdkafka](https://github.com/edenhill/librdkafka) можно включить через конфигурацию.

```xml theme={null}
<kafka>
   <debug>all</debug>
</kafka>
```

<div id="handling-malformed-messages">
  ##### Обработка некорректных сообщений
</div>

Kafka часто используют как «свалку» для данных. Из-за этого в топиках оказываются сообщения в разных форматах и с несогласованными именами полей. По возможности избегайте этого и используйте возможности Kafka, такие как Kafka Streams или ksqlDB, чтобы сообщения были корректно сформированы и согласованы еще до записи в Kafka. Если это невозможно, в ClickHouse есть несколько возможностей, которые могут помочь.

* Обрабатывайте поле сообщения как строки. При необходимости в операторе materialized view можно использовать функции для очистки и приведения типов. Это не решение для production, но может помочь при разовой ингестии.
* Если вы читаете JSON из топика в формате JSONEachRow, используйте настройку [`input_format_skip_unknown_fields`](/docs/ru/reference/settings/formats#input_format_skip_unknown_fields). При записи данных ClickHouse по умолчанию генерирует исключение, если входные данные содержат столбцы, которых нет в целевой таблице. Однако если эта опция включена, такие лишние столбцы будут игнорироваться. Опять же, это не решение уровня production и оно может запутать других.
* Обратите внимание на настройку `kafka_skip_broken_messages`. Она требует, чтобы пользователь задал допустимый уровень ошибок для некорректных сообщений на block с учетом `kafka_max_block_size`. Если этот порог превышен (в абсолютном количестве сообщений), снова будет применяться стандартное поведение с исключением, а остальные сообщения будут пропущены.

<div id="delivery-semantics-and-challenges-with-duplicates">
  ##### Семантика доставки и проблемы с дубликатами
</div>

Kafka движок таблицы имеет семантику at-least-once. Дубликаты возможны в ряде известных, хотя и редких, случаев. Например, сообщения могут быть прочитаны из Kafka и успешно вставлены в ClickHouse. Но до того, как новое смещение будет зафиксировано, соединение с Kafka может быть потеряно. В такой ситуации требуется повторная попытка обработки блока. Блок может быть [дедуплицирован](/docs/ru/reference/engines/table-engines/mergetree-family/replication) при использовании distributed таблицы или ReplicatedMergeTree в качестве целевой таблицы. Хотя это снижает вероятность появления дублирующихся строк, такой подход опирается на идентичность блоков. Такие события, как ребалансировка Kafka, могут нарушить это допущение, что в редких случаях приводит к дубликатам.

<div id="quorum-based-inserts">
  ##### Вставки на основе кворума
</div>

В некоторых случаях, когда в ClickHouse нужны более высокие гарантии доставки, могут потребоваться [вставки на основе кворума](/docs/ru/reference/settings/session-settings#insert_quorum). Это нельзя настроить для materialized view или целевой таблицы. Однако это можно задать для пользовательских профилей, например.

```xml theme={null}
<profiles>
  <default>
    <insert_quorum>2</insert_quorum>
  </default>
</profiles>
```

<div id="clickhouse-to-kafka">
  ### Из ClickHouse в Kafka
</div>

Хотя это и более редкий сценарий использования, данные ClickHouse также можно сохранять в Kafka. Например, мы вручную вставим строки в таблицу с движком таблицы Kafka. Эти данные будут прочитаны тем же Kafka-движком, а его materialized view поместит их в таблицу MergeTree. Наконец, мы покажем, как использовать materialized views при вставке в Kafka, чтобы наполнять таблицы на основе существующих исходных таблиц.

<div id="steps-1">
  #### Шаги
</div>

Нашу первоначальную цель лучше всего иллюстрирует следующая схема:

<Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/kafka_02.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=9093dc39ca712ae891de358b107bcfad" size="lg" alt="Схема движка таблицы Kafka со вставками" width="2048" height="852" data-path="images/integrations/data-ingestion/kafka/kafka_02.webp" />

Мы предполагаем, что вы уже создали таблицы и представления, описанные в разделе [Kafka в ClickHouse](#kafka-to-clickhouse), и что топик был полностью прочитан.

<Steps>
  <Step title="Прямая вставка строк" id="1-inserting-rows-directly">
    Сначала проверьте количество строк в целевой таблице.

    ```sql theme={null}
    SELECT count() FROM github;
    ```

    У вас должно получиться 200 000 строк:

    ```response theme={null}
    ┌─count()─┐
    │  200000 │
    └─────────┘
    ```

    Теперь выполните вставку строк из целевой таблицы GitHub обратно в движок таблицы Kafka github\_queue. Обратите внимание, что мы используем формат JSONEachRow и ограничиваем выборку до 100 строк с помощью LIMIT.

    ```sql theme={null}
    INSERT INTO github_queue SELECT * FROM github LIMIT 100 FORMAT JSONEachRow
    ```

    Пересчитайте строки в GitHub, чтобы убедиться, что их количество увеличилось на 100. Как показано на приведённой выше диаграмме, строки были вставлены в Kafka через движок таблицы Kafka, после чего тот же движок снова их прочитал и вставил в целевую таблицу GitHub с помощью нашего materialized view!

    ```sql theme={null}
    SELECT count() FROM github;
    ```

    Вы увидите ещё 100 строк:

    ```response theme={null}
    ┌─count()─┐
    │  200100 │
    └─────────┘
    ```
  </Step>

  <Step title="Использование materialized views" id="2-using-materialized-views">
    Мы можем использовать materialized views, чтобы отправлять сообщения в движок Kafka (и в топик), когда документы вставляются в таблицу. Когда строки вставляются в таблицу GitHub, срабатывает materialized view, в результате чего строки снова вставляются в движок Kafka и в новый топик. Это лучше всего показано на иллюстрации:

    <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/kafka_03.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=ead2a8c956700436a1fd15bb803ef49b" size="lg" alt="Диаграмма движка таблицы Kafka с materialized views" width="2048" height="870" data-path="images/integrations/data-ingestion/kafka/kafka_03.webp" />

    Создайте новый топик Kafka `github_out` или его эквивалент. Убедитесь, что движок таблицы Kafka `github_out_queue` указывает на этот топик.

    ```sql theme={null}
    CREATE TABLE github_out_queue
    (
        file_time DateTime,
        event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4, 'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
        actor_login LowCardinality(String),
        repo_name LowCardinality(String),
        created_at DateTime,
        updated_at DateTime,
        action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
        comment_id UInt64,
        path String,
        ref LowCardinality(String),
        ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
        creator_user_login LowCardinality(String),
        number UInt32,
        title String,
        labels Array(LowCardinality(String)),
        state Enum('none' = 0, 'open' = 1, 'closed' = 2),
        assignee LowCardinality(String),
        assignees Array(LowCardinality(String)),
        closed_at DateTime,
        merged_at DateTime,
        merge_commit_sha String,
        requested_reviewers Array(LowCardinality(String)),
        merged_by LowCardinality(String),
        review_comments UInt32,
        member_login LowCardinality(String)
    )
       ENGINE = Kafka('host:port', 'github_out', 'clickhouse_out',
                'JSONEachRow') SETTINGS kafka_thread_per_consumer = 0, kafka_num_consumers = 1;
    ```

    Теперь создайте новое materialized view `github_out_mv`, которое будет указывать на таблицу GitHub и при срабатывании выполнять вставку строк в указанный выше движок. В результате новые записи из таблицы GitHub будут отправляться в наш новый топик Kafka.

    ```sql theme={null}
    CREATE MATERIALIZED VIEW github_out_mv TO github_out_queue AS
    SELECT file_time, event_type, actor_login, repo_name,
           created_at, updated_at, action, comment_id, path,
           ref, ref_type, creator_user_login, number, title,
           labels, state, assignee, assignees, closed_at, merged_at,
           merge_commit_sha, requested_reviewers, merged_by,
           review_comments, member_login
    FROM github
    FORMAT JsonEachRow;
    ```

    Если выполнить вставку в исходный топик github, созданный в разделе [Kafka в ClickHouse](#kafka-to-clickhouse), документы автоматически появятся в топике "github\_clickhouse". Убедитесь в этом с помощью встроенных инструментов Kafka. Например, ниже мы выполняем вставку 100 строк в топик github с помощью [kcat](https://github.com/edenhill/kcat) для топика, размещённого в Confluent Cloud:

    ```sql theme={null}
    head -n 10 github_all_columns.ndjson |
    kcat -P \
      -b <host>:<port> \
      -t github
      -X security.protocol=sasl_ssl \
      -X sasl.mechanisms=PLAIN \
      -X sasl.username=<username> \
      -X sasl.password=<password>
    ```

    Чтение из топика `github_out` должно подтвердить доставку сообщений.

    ```sql theme={null}
    kcat -C \
      -b <host>:<port> \
      -t github_out \
      -X security.protocol=sasl_ssl \
      -X sasl.mechanisms=PLAIN \
      -X sasl.username=<username> \
      -X sasl.password=<password> \
      -e -q |
    wc -l
    ```

    Несмотря на то что пример достаточно сложный, он наглядно демонстрирует возможности materialized views при использовании совместно с движком Kafka.
  </Step>
</Steps>

<div id="clusters-and-performance">
  ### Кластеры и производительность
</div>

<div id="working-with-clickhouse-clusters">
  #### Работа с кластерами ClickHouse
</div>

Благодаря группам потребителей Kafka несколько экземпляров ClickHouse могут читать из одного и того же топика. Каждому потребителю назначается одна партиция топика в соотношении 1:1. При масштабировании чтения ClickHouse с использованием движка таблицы Kafka учитывайте, что общее число потребителей в кластере не может превышать число партиций в топике. Поэтому заранее убедитесь, что для топика настроено достаточное количество партиций.

Несколько экземпляров ClickHouse можно настроить на чтение из топика с одним и тем же идентификатором группы потребителей, указанным при создании движка таблицы Kafka. Таким образом, каждый экземпляр будет читать из одной или нескольких партиций, записывая сегменты в свою локальную целевую таблицу. Целевые таблицы, в свою очередь, можно настроить на использование ReplicatedMergeTree для устранения дублирования данных. Такой подход позволяет масштабировать чтение из Kafka вместе с кластером ClickHouse при условии, что в Kafka достаточно партиций.

<Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/kafka_04.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=fc14aa743c7b5cda9b7d722e18bc59f0" size="lg" alt="Диаграмма движка таблицы Kafka с кластерами ClickHouse" width="2048" height="1094" data-path="images/integrations/data-ingestion/kafka/kafka_04.webp" />

<div id="tuning-performance">
  #### Настройка производительности
</div>

При увеличении пропускной способности таблицы с движком Kafka учитывайте следующее:

* Производительность зависит от размера сообщений, формата и типов целевых таблиц. Показатель в 100 тыс. строк/с для одного движка таблицы считается достижимым. По умолчанию сообщения читаются блоками; это регулируется параметром kafka\_max\_block\_size. По умолчанию он равен [max\_insert\_block\_size](/docs/ru/reference/settings/session-settings#max_insert_block_size), то есть 1,048,576. Если сообщения не являются исключительно большими, это значение почти всегда стоит увеличивать. Значения в диапазоне от 500 тыс. до 1 млн вполне обычны. Протестируйте и оцените, как это влияет на пропускную способность.
* Число потребителей для движка таблицы можно увеличить с помощью kafka\_num\_consumers. Однако по умолчанию вставки будут выполняться последовательно в одном потоке, если kafka\_thread\_per\_consumer не изменить относительно значения по умолчанию 1. Установите это значение в 1, чтобы сбросы на диск выполнялись параллельно. Обратите внимание: создание таблицы с движком Kafka с N потребителями (и kafka\_thread\_per\_consumer=1) логически эквивалентно созданию N движков Kafka, каждый со своей materialized view и kafka\_thread\_per\_consumer=0.
* Увеличение числа потребителей не даётся бесплатно. Каждый потребитель поддерживает собственные буферы и потоки, увеличивая накладные расходы на сервер. Учитывайте эти накладные расходы и по возможности сначала линейно масштабируйте нагрузку по кластеру.
* Если пропускная способность потока сообщений Kafka непостоянна, а задержки допустимы, рассмотрите увеличение stream\_flush\_interval\_ms, чтобы на диск сбрасывались более крупные блоки.
* [background\_message\_broker\_schedule\_pool\_size](/docs/ru/reference/settings/server-settings/settings#background_message_broker_schedule_pool_size) задаёт количество потоков, выполняющих фоновые задачи. Эти потоки используются для стриминга Kafka. Этот параметр применяется при запуске сервера ClickHouse и не может быть изменён в пользовательском сеансе; по умолчанию его значение равно 16. Если вы видите тайм-ауты в журнале, возможно, стоит увеличить это значение.
* Для взаимодействия с Kafka используется библиотека librdkafka, которая также создаёт потоки. Поэтому большое число таблиц Kafka или потребителей может приводить к большому количеству переключений контекста. Либо распределите эту нагрузку по кластеру, по возможности реплицируя только целевые таблицы, либо рассмотрите использование движка таблицы для чтения из нескольких топиков — поддерживается список значений. Из одной таблицы можно читать через несколько materialized view, каждая из которых фильтрует данные из определённого топика.

Любые изменения настроек следует тестировать. Мы рекомендуем отслеживать отставание потребителей Kafka, чтобы убедиться, что масштабирование выбрано правильно.

<div id="additional-settings">
  #### Дополнительные настройки
</div>

Помимо настроек, рассмотренных выше, интерес могут представлять следующие:

* [Kafka\_max\_wait\_ms](/docs/ru/reference/settings/session-settings#kafka_max_wait_ms) — Время ожидания в миллисекундах при чтении сообщений из Kafka перед повторной попыткой. Задаётся на уровне профиля пользователя; значение по умолчанию — 5000.

[Все настройки ](https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md)базовой библиотеки librdkafka также можно указывать в файлах конфигурации ClickHouse внутри элемента *kafka* — имена настроек должны быть XML-элементами, в которых точки заменены на подчёркивания, например.

```xml theme={null}
<clickhouse>
   <kafka>
       <enable_ssl_certificate_verification>false</enable_ssl_certificate_verification>
   </kafka>
</clickhouse>
```

Это настройки для опытных пользователей, и мы рекомендуем обратиться к документации Kafka за более подробными разъяснениями.
