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

> Usando o motor de tabela Kafka

# Usando o motor de tabela 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>;
};

O motor de tabela Kafka pode ser usado para [**ler** dados do](#kafka-to-clickhouse) e [**gravar** dados no](#clickhouse-to-kafka) Apache Kafka e em outros brokers compatíveis com a API do Kafka (por exemplo, Redpanda, Amazon MSK).

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

<Note>
  Se você usa ClickHouse Cloud, recomendamos usar o [ClickPipes](/docs/pt-BR/integrations/clickpipes/home). O ClickPipes oferece suporte nativo a conexões de rede privadas, ao dimensionamento independente da ingestão e dos recursos do cluster, além de monitoramento abrangente para streaming de dados do Kafka para o ClickHouse.
</Note>

Para usar o mecanismo de tabela Kafka, você deve ter familiaridade geral com [visões materializadas do ClickHouse](/docs/pt-BR/concepts/features/materialized-views/cascading-materialized-views).

<div id="overview">
  #### Visão geral
</div>

Inicialmente, focamos no caso de uso mais comum: usar o motor de tabela Kafka para inserir dados do Kafka no ClickHouse.

O motor de tabela Kafka permite que o ClickHouse leia diretamente de um tópico Kafka. Embora seja útil para visualizar mensagens em um tópico, esse motor, por definição, permite apenas uma leitura única; ou seja, quando uma consulta é executada na tabela, ele consome dados da fila e avança o offset do consumer antes de retornar os resultados ao chamador. Na prática, os dados não podem ser lidos novamente sem redefinir esses offsets.

Para persistir esses dados a partir de uma leitura do motor de tabela, precisamos de uma forma de capturá-los e inseri-los em outra tabela. Visões materializadas baseadas em trigger fornecem essa funcionalidade de forma nativa. Uma visão materializada inicia uma leitura no motor de tabela, recebendo lotes de documentos. A cláusula TO determina o destino dos dados — normalmente uma tabela da [família MergeTree](/docs/pt-BR/reference/engines/table-engines/mergetree-family/index). Esse processo é ilustrado abaixo:

<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="Diagrama da arquitetura do motor de tabela Kafka" style={{width: '80%'}} width="2048" height="852" data-path="images/integrations/data-ingestion/kafka/kafka_01.webp" />

<div id="steps">
  #### Etapas
</div>

<Steps>
  <Step title="Prepare" id="1-prepare">
    Se você tiver dados carregados em um tópico de destino, poderá adaptar o conteúdo a seguir para uso no seu conjunto de dados. Como alternativa, um conjunto de dados do GitHub de exemplo é fornecido [aqui](https://datasets-documentation.s3.eu-west-3.amazonaws.com/kafka/github_all_columns.ndjson). Esse conjunto de dados é usado nos exemplos abaixo e utiliza um esquema reduzido e um subconjunto das linhas (especificamente, limitamos aos eventos do GitHub relacionados ao [repositório do ClickHouse](https://github.com/ClickHouse/ClickHouse)), em comparação com o conjunto de dados completo disponível [aqui](https://ghe.clickhouse.tech/), por questão de brevidade. Ainda assim, isso é suficiente para que a maioria das consultas [publicadas com o conjunto de dados](https://ghe.clickhouse.tech/) funcione.
  </Step>

  <Step title="Configure o ClickHouse" id="2-configure-clickhouse">
    Esta etapa é necessária se você estiver se conectando a um Kafka com segurança habilitada. Essas configurações não podem ser passadas por comandos SQL DDL e devem ser configuradas no config.xml do ClickHouse. Presumimos que você esteja se conectando a uma instância protegida por SASL. Este é o método mais simples ao interagir com a 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>
    ```

    Coloque o trecho acima em um novo arquivo no diretório `conf.d/` ou mescle-o aos arquivos de configuração existentes. Para as configurações que podem ser definidas, consulte [aqui](/docs/pt-BR/reference/engines/table-engines/integrations/kafka#configuration).

    Também vamos criar um banco de dados chamado `KafkaEngine` para usar neste tutorial:

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

    Depois de criar o banco de dados, você precisará mudar para ele:

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

  <Step title="Crie a tabela de destino" id="3-create-the-destination-table">
    Prepare a sua tabela de destino. No exemplo abaixo, usamos o esquema reduzido do GitHub por questão de brevidade. Observe que, embora usemos um motor de tabela MergeTree, este exemplo pode ser facilmente adaptado para qualquer membro da [família MergeTree](/docs/pt-BR/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="Crie e popule o tópico" id="4-create-and-populate-the-topic">
    Em seguida, vamos criar um tópico. Há várias ferramentas que podemos usar para fazer isso. Se estivermos executando o Kafka localmente em nossa máquina ou em um contêiner Docker, o [RPK](https://docs.redpanda.com/current/get-started/rpk-install/) funciona bem. Podemos criar um tópico chamado `github` com 5 partições executando o seguinte comando:

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

    Se estivermos executando o Kafka no Confluent Cloud, talvez prefiramos usar a [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
    ```

    Agora precisamos popular este tópico com alguns dados, o que faremos usando [kcat](https://github.com/edenhill/kcat). Podemos executar um comando semelhante ao seguinte se estivermos executando o Kafka localmente com a autenticação desabilitada:

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

    Ou o seguinte, se nosso cluster Kafka usar SASL para autenticação:

    ```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> \
    ```

    O conjunto de dados contém 200.000 linhas, portanto a ingestão deve levar apenas alguns segundos. Se você quiser trabalhar com um conjunto de dados maior, dê uma olhada [na seção sobre grandes conjuntos de dados](https://github.com/ClickHouse/kafka-samples/tree/main/producer#large-datasets) do repositório do GitHub [ClickHouse/kafka-samples](https://github.com/ClickHouse/kafka-samples).
  </Step>

  <Step title="Crie o motor de tabela Kafka" id="5-create-the-kafka-table-engine">
    O exemplo abaixo cria um motor de tabela com o mesmo esquema da tabela MergeTree. Isso não é estritamente necessário, pois você pode ter colunas alias ou efêmeras na tabela de destino. As configurações são importantes; no entanto, observe o uso de `JSONEachRow` como tipo de dado para consumir JSON de um tópico Kafka. Os valores `github` e `clickhouse` representam, respectivamente, o nome do tópico e o nome do grupo de consumidores. Na verdade, os tópicos podem ser especificados como uma lista de valores.

    ```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;
    ```

    Discutimos abaixo as configurações do motor e o ajuste de desempenho. Neste ponto, um `select` simples na tabela `github_queue` deve ler algumas linhas. Observe que isso avançará os offsets do consumidor, impedindo que essas linhas sejam lidas novamente sem um [reset](#common-operations). Observe o limite e o parâmetro obrigatório `stream_like_engine_allow_direct_select.`
  </Step>

  <Step title="Crie a visão materializada" id="6-create-the-materialized-view">
    A visão materializada conectará as duas tabelas criadas anteriormente, lendo dados do motor de tabela Kafka e inserindo-os na tabela MergeTree de destino. Podemos fazer várias transformações nos dados. Faremos uma leitura e inserção simples. O uso de \* pressupõe que os nomes das colunas sejam idênticos (diferenciam maiúsculas de minúsculas).

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

    No momento da criação, a visão materializada se conecta ao motor Kafka e começa a ler, inserindo linhas na tabela de destino. Esse processo continuará indefinidamente, consumindo as mensagens inseridas posteriormente no Kafka. Sinta-se à vontade para executar novamente o script de inserção para inserir mais mensagens no Kafka.
  </Step>

  <Step title="Confirme se as linhas foram inseridas" id="7-confirm-rows-have-been-inserted">
    Confirme que existem dados na tabela de destino:

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

    Você deve ver 200.000 linhas:

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

<div id="common-operations">
  #### Operações comuns
</div>

<div id="stopping--restarting-message-consumption">
  ##### Interrompendo & reiniciando o consumo de mensagens
</div>

Para interromper o consumo de mensagens, você pode desanexar a tabela com motor Kafka:

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

Isso não afetará os offsets do grupo de consumidores. Para reiniciar o consumo e continuar a partir do offset anterior, anexe a tabela novamente.

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

<div id="adding-kafka-metadata">
  ##### Adicionando metadados do Kafka
</div>

Pode ser útil acompanhar os metadados das mensagens originais do Kafka após a ingestão no ClickHouse. Por exemplo, talvez queiramos saber quanto de um tópico ou partição específicos já consumimos. Para isso, o motor de tabela Kafka expõe várias [colunas virtuais](/docs/pt-BR/reference/engines/table-engines/index#table_engines-virtual_columns). Elas podem ser persistidas como colunas na nossa tabela de destino, modificando o esquema e a instrução `SELECT` da visão materializada.

Primeiro, executamos a operação de parada descrita acima antes de adicionar colunas à nossa tabela de destino.

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

Abaixo, adicionamos colunas informativas para identificar o tópico de origem e a partição de onde a linha se originou.

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

Em seguida, precisamos garantir que as colunas virtuais sejam mapeadas conforme necessário.
As colunas virtuais são prefixadas com `_`.
Uma lista completa das colunas virtuais pode ser encontrada [aqui](/docs/pt-BR/reference/engines/table-engines/integrations/kafka#virtual-columns).

Para atualizar nossa tabela com as colunas virtuais, precisaremos remover a visão materializada, reanexar a tabela do motor Kafka e recriar a visão materializada.

```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;
```

As linhas recém-consumidas devem conter os metadados.

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

O resultado fica assim:

| 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">
  ##### Modificar as configurações do motor Kafka
</div>

Recomendamos excluir a tabela do motor Kafka e recriá-la com as novas configurações. A visão materializada não precisa ser modificada durante esse processo - o consumo de mensagens será retomado assim que a tabela do motor Kafka for recriada.

<div id="debugging-issues">
  ##### Depuração de problemas
</div>

Erros como falhas de autenticação não são reportados nas respostas aos comandos DDL do motor Kafka. Para diagnosticar esses problemas, recomendamos usar o principal arquivo de log do ClickHouse, clickhouse-server.err.log. Também é possível habilitar logs de rastreamento adicionais para a biblioteca cliente Kafka subjacente [librdkafka](https://github.com/edenhill/librdkafka) por meio da configuração.

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

<div id="handling-malformed-messages">
  ##### Tratando mensagens malformadas
</div>

O Kafka costuma ser usado como um "depósito" de dados. Isso faz com que os tópicos contenham formatos de mensagem mistos e nomes de campo inconsistentes. Evite isso e use recursos do Kafka, como Kafka Streams ou ksqlDB, para garantir que as mensagens estejam bem formadas e consistentes antes da inserção no Kafka. Se essas opções não forem viáveis, o ClickHouse oferece alguns recursos que podem ajudar.

* Trate os campos da mensagem como strings. É possível usar funções na instrução da visão materializada para fazer limpeza e conversão de tipo, se necessário. Isso não deve ser considerado uma solução de produção, mas pode ajudar em uma ingestão pontual.
* Se você estiver consumindo JSON de um tópico usando o formato JSONEachRow, use a configuração [`input_format_skip_unknown_fields`](/docs/pt-BR/reference/settings/formats#input_format_skip_unknown_fields). Ao gravar dados, por padrão, o ClickHouse lança uma exceção se os dados de entrada contiverem colunas que não existem na tabela de destino. No entanto, se essa opção estiver habilitada, essas colunas excedentes serão ignoradas. Novamente, isso não é uma solução adequada para produção e pode confundir outras pessoas.
* Considere a configuração `kafka_skip_broken_messages`. Isso exige que o usuário especifique o nível de tolerância por bloco para mensagens malformadas, considerado no contexto de `kafka_max_block_size`. Se essa tolerância for excedida (medida em número absoluto de mensagens), o comportamento normal de exceção será retomado, e as demais mensagens serão ignoradas.

<div id="delivery-semantics-and-challenges-with-duplicates">
  ##### Semântica de entrega e desafios com duplicatas
</div>

O motor de tabela Kafka tem semântica de entrega de pelo menos uma vez. Duplicatas podem ocorrer em várias circunstâncias raras já conhecidas. Por exemplo, as mensagens podem ser lidas do Kafka e inseridas com sucesso no ClickHouse. Antes que o novo offset possa ser confirmado, a conexão com o Kafka pode ser perdida. Nessa situação, é necessário tentar novamente o bloco. O bloco pode ser [deduplicado ](/docs/pt-BR/reference/engines/table-engines/mergetree-family/replication)usando uma tabela distribuída ou ReplicatedMergeTree como tabela de destino. Embora isso reduza a chance de linhas duplicadas, depende de blocos idênticos. Eventos como um rebalanceamento do Kafka podem invalidar essa premissa, causando duplicatas em circunstâncias raras.

<div id="quorum-based-inserts">
  ##### Inserções com quórum
</div>

Você pode precisar de [inserções com quórum](/docs/pt-BR/reference/settings/session-settings#insert_quorum) nos casos em que sejam necessárias garantias de entrega mais fortes no ClickHouse. Isso não pode ser configurado na visão materializada nem na tabela de destino. No entanto, pode ser configurado para perfis de usuário, por exemplo.

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

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

Embora seja um caso de uso menos comum, os dados do ClickHouse também podem ser persistidos no Kafka. Por exemplo, vamos inserir linhas manualmente em uma tabela com motor Kafka. Esses dados serão lidos pelo mesmo motor Kafka, cuja visão materializada gravará os dados em uma tabela Merge Tree. Por fim, demonstramos a aplicação de visões materializadas em inserções no Kafka para ler dados de tabelas de origem existentes.

<div id="steps-1">
  #### Etapas
</div>

Nosso objetivo inicial é ilustrado da melhor forma a seguir:

<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="Diagrama do motor de tabela Kafka com inserções" width="2048" height="852" data-path="images/integrations/data-ingestion/kafka/kafka_02.webp" />

Pressupomos que você tenha as tabelas e views criadas nas etapas de [Kafka para ClickHouse](#kafka-to-clickhouse) e que o tópico tenha sido totalmente consumido.

<Steps>
  <Step title="Inserção direta de linhas" id="1-inserting-rows-directly">
    Primeiro, confirme o número de linhas da tabela de destino.

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

    Você deve ter 200.000 linhas:

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

    Agora insira linhas da tabela de destino GitHub de volta para o motor de tabela Kafka github\_queue. Observe como utilizamos o formato JSONEachRow e limitamos o select a 100.

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

    Conte novamente as linhas na tabela GitHub para confirmar que o total aumentou em 100. Como mostrado no diagrama acima, as linhas foram inseridas no Kafka por meio do motor de tabela Kafka antes de serem relidas pelo mesmo motor e inseridas na tabela de destino GitHub pela nossa visão materializada!

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

    Você deve ver mais 100 linhas:

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

  <Step title="Usando visões materializadas" id="2-using-materialized-views">
    Podemos usar visões materializadas para enviar mensagens a um motor Kafka (e a um tópico) quando documentos são inseridos em uma tabela. Quando linhas são inseridas na tabela GitHub, uma visão materializada é acionada, fazendo com que as linhas sejam reinseridas em um motor Kafka e em um novo tópico. Mais uma vez, isso fica mais claro na ilustração:

    <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="Diagrama do motor de tabela Kafka com visões materializadas" width="2048" height="870" data-path="images/integrations/data-ingestion/kafka/kafka_03.webp" />

    Crie um novo tópico Kafka `github_out` ou equivalente. Certifique-se de que um motor de tabela Kafka `github_out_queue` aponte para esse tópico.

    ```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;
    ```

    Agora crie uma nova visão materializada `github_out_mv` que aponte para a tabela GitHub, inserindo linhas no mecanismo acima quando ela for acionada. Como resultado, as adições à tabela GitHub serão enviadas ao nosso novo tópico 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;
    ```

    Se você inserir no tópico github original, criado como parte de [Kafka to ClickHouse](#kafka-to-clickhouse), os documentos aparecerão como mágica no tópico "github\_clickhouse". Confirme isso com ferramentas nativas do Kafka. Por exemplo, abaixo, inserimos 100 linhas no tópico github usando [kcat](https://github.com/edenhill/kcat) em um tópico hospedado no 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>
    ```

    Uma leitura no tópico `github_out` deve confirmar a entrega das mensagens.

    ```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
    ```

    Embora seja um exemplo mais elaborado, isso ilustra o poder das visões materializadas quando usadas em conjunto com o motor Kafka.
  </Step>
</Steps>

<div id="clusters-and-performance">
  ### Clusters e desempenho
</div>

<div id="working-with-clickhouse-clusters">
  #### Trabalhando com clusters do ClickHouse
</div>

Por meio de grupos de consumidores do Kafka, várias instâncias do ClickHouse podem ler do mesmo tópico. Cada consumidor será atribuído a uma partição do tópico em um mapeamento 1:1. Ao escalar o consumo do ClickHouse usando o motor de tabela Kafka, considere que o número total de consumidores em um cluster não pode exceder o número de partições no tópico. Portanto, garanta antecipadamente que o particionamento esteja configurado adequadamente para o tópico.

Várias instâncias do ClickHouse podem ser configuradas para ler de um tópico usando o mesmo ID de grupo de consumidores, especificado durante a criação do motor de tabela Kafka. Assim, cada instância lerá de uma ou mais partições, inserindo segmentos em sua tabela de destino local. As tabelas de destino, por sua vez, podem ser configuradas para usar um ReplicatedMergeTree para lidar com a duplicação dos dados. Essa abordagem permite escalar as leituras do Kafka com o cluster do ClickHouse, desde que haja partições do Kafka suficientes.

<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="Diagrama do motor de tabela Kafka com clusters do ClickHouse" width="2048" height="1094" data-path="images/integrations/data-ingestion/kafka/kafka_04.webp" />

<div id="tuning-performance">
  #### Ajuste de desempenho
</div>

Considere os pontos a seguir para aumentar a taxa de transferência de tabelas com motor Kafka:

* O desempenho varia conforme o tamanho da mensagem, o formato e os tipos da tabela de destino. É razoável esperar 100 mil linhas/s em um único motor de tabela. Por padrão, as mensagens são lidas em blocos, controlados pelo parâmetro `kafka_max_block_size`. Por padrão, ele é definido como [max\_insert\_block\_size](/docs/pt-BR/reference/settings/session-settings#max_insert_block_size), cujo valor padrão é 1.048.576. A menos que as mensagens sejam extremamente grandes, quase sempre vale a pena aumentar esse valor. Valores entre 500 mil e 1 milhão não são incomuns. Teste e avalie o efeito sobre a taxa de transferência.
* O número de consumers de um motor de tabela pode ser aumentado com `kafka_num_consumers`. No entanto, por padrão, os inserts serão linearizados em uma única thread, a menos que `kafka_thread_per_consumer` seja alterado do valor padrão 1. Defina-o como 1 para garantir que os flushes sejam realizados em paralelo. Observe que criar uma tabela com motor Kafka com N consumers (e `kafka_thread_per_consumer=1`) é logicamente equivalente a criar N motores Kafka, cada um com uma visão materializada e `kafka_thread_per_consumer=0`.
* Aumentar o número de consumers não é gratuito. Cada consumer mantém seus próprios buffers e threads, aumentando a sobrecarga no servidor. Fique atento a essa sobrecarga e, se possível, primeiro faça o scale linearmente no cluster.
* Se a taxa de transferência das mensagens do Kafka variar e atrasos forem aceitáveis, considere aumentar `stream_flush_interval_ms` para garantir que blocos maiores sejam gravados.
* [background\_message\_broker\_schedule\_pool\_size](/docs/pt-BR/reference/settings/server-settings/settings#background_message_broker_schedule_pool_size) define o número de threads que executam tasks em segundo plano. Essas threads são usadas para streaming do Kafka. Essa configuração é aplicada na inicialização do servidor ClickHouse e não pode ser alterada em uma sessão de usuário; o valor padrão é 16. Se você observar timeouts nos logs, pode ser apropriado aumentá-la.
* Para a comunicação com o Kafka, é usada a biblioteca librdkafka, que por sua vez cria threads. Assim, um grande número de tabelas Kafka ou de consumers pode resultar em muitas trocas de contexto. Distribua essa carga pelo cluster, replicando apenas as tabelas de destino, se possível, ou considere usar um motor de tabela para ler de vários topics — uma lista de valores é compatível. Várias visões materializadas podem ler de uma única tabela, cada uma filtrando os dados de um topic específico.

Quaisquer alterações de configuração devem ser testadas. Recomendamos monitorar a defasagem dos consumers do Kafka para garantir que o dimensionamento esteja adequado.

<div id="additional-settings">
  #### Configurações adicionais
</div>

Além das configurações discutidas acima, as seguintes podem ser de interesse:

* [Kafka\_max\_wait\_ms](/docs/pt-BR/reference/settings/session-settings#kafka_max_wait_ms) - O tempo de espera, em milissegundos, para ler mensagens do Kafka antes de tentar novamente. É definido no nível do perfil do usuário e o padrão é 5000.

[Todas as configurações ](https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md)da biblioteca librdkafka subjacente também podem ser colocadas nos arquivos de configuração do ClickHouse dentro de um elemento *kafka* - os nomes das configurações devem ser elementos XML com os pontos substituídos por sublinhados, por exemplo.

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

Estas são configurações avançadas, e sugerimos que você consulte a documentação do Kafka para uma explicação mais detalhada.
