Pular para o conteúdo principal

Input Kafka

Consome mensagens de um Kafka que já existe — o da própria plataforma ou um cluster externo — e as traz para dentro da pipeline.

É o Input usado quando o dado já está em Kafka e você quer processá-lo ou entregá-lo em outro lugar.

Campos

CampoO que éPadrão
TópicosUm ou mais tópicos de origem, separados por vírgula. Ex.: pedidos, pedidos-cancelados.
Group IDO marcador de posição de leitura. É ele que guarda até onde a pipeline já leu. Dois consumidores com o mesmo Group ID dividem as partições entre si; com Group IDs diferentes, cada um lê tudo.
Conexão Kafka (padrão)Qual cluster consumir. Deixe em Padrão do pipeline para usar o broker da própria pipeline, ou escolha uma conexão cadastrada em Kafka → Conexões para consumir de um cluster externo.Padrão do pipeline

A conexão escolhida aqui é o padrão — ela pode ser sobrescrita por cluster no momento do deploy. Veja Conexões Kafka.

Vários tópicos não viram um só

Consumir de dois tópicos não junta os dois numa fila única. Cada tópico de entrada continua separado ao longo da pipeline: o Input cria um tópico de saída para cada um, e você pode tratá-los de formas diferentes adiante.

Se você quisesse mesmo juntar tudo, ligaria os dois caminhos ao mesmo Output.

O que este Input escreve no tópico

Este Input repassa a mensagem inteira, sem alterar nada:

ParteConteúdo
ValueO mesmo Value que veio do tópico de origem
KeyA mesma Key que veio
HeadersOs mesmos headers que vieram

Exemplo

Mensagem no tópico de origem pedidos:

Key: {"pedido":1234}
Value: {"pedido":1234,"cliente":"ACME","total":199.90}
Headers: origem=erp versao=2

Mensagem escrita no tópico de saída da pipeline — idêntica:

Key: {"pedido":1234}
Value: {"pedido":1234,"cliente":"ACME","total":199.90}
Headers: origem=erp versao=2

Essa fidelidade é o que permite encadear pipelines: se o tópico de origem já carrega os headers que o SQL Server CDC produz, eles chegam intactos ao Output de PostgreSQL ou Snowflake, e a replicação funciona igual.

Combinações comuns

QueroMonte assim
Levar um fluxo Kafka existente para o ElasticsearchKafka → Output Elasticsearch
Arquivar um tópico no data lakeKafka → Output GCS
Filtrar ou reformatar antes de entregarKafka → Transform → Output
Copiar de um cluster para outroKafka (conexão A) → Output Kafka (conexão B)