Mover dados usando um conector Kafka

Concluído

Apache Kafka é uma plataforma de software livre usada para transmitir eventos de maneira distribuída. Muitas empresas usam o Kafka para cenários de integração de dados de alto desempenho em grande escala. O Kafka Connect é uma ferramenta em seu pacote para transmitir dados entre o Kafka e outros sistemas de dados. De maneira compreensível, isso pode incluir o Azure Cosmos DB como uma fonte de dados ou um destino (coletor) de dados.

Instalação

Os conectores do Kafka Connect para o Azure Cosmos DB estão disponíveis como um projeto de código aberto no GitHub em microsoft/kafka-connect-cosmosdb. As instruções para baixar e instalar o arquivo JAR manualmente estão disponíveis no repositório.

Configuração

Quatro propriedades de configuração devem ser definidas para configurar corretamente a conectividade com uma conta do Azure Cosmos DB for NoSQL.

Propriedade Valor
connect.cosmos.connection.endpoint URI do ponto de extremidade da conta
connect.cosmos.master.key Chave de conta
connect.cosmos.databasename Nome do recurso do banco de dados
connect.cosmos.containers.topicmap Usando o formato CSV, um mapeamento dos tópicos do Kafka para contêineres

Mapa de tópicos para contêineres

Cada contêiner deve ser mapeado para um tópico. Por exemplo, suponha que você gostaria que o contêiner de produtos fosse mapeado para o tópico prodlistener e o contêiner de clientes para o tópico custlistener. Nesse caso, você deve usar a seguinte cadeia de caracteres de mapeamento CSV: prodlistener#products,custlistener#customers.

Gravar no Azure Cosmos DB

Vamos gravar dados no Azure Cosmos DB criando um tópico. No Apache Kafka, todas as mensagens são enviadas por meio de tópicos.

Você pode criar um tópico usando o comando kafka-topics. Este exemplo criará um tópico chamado prodlistener.

kafka-topics --create \
    --zookeeper localhost:2181 \
    --topic prodlistener \
    --replication-factor 1 \
    --partitions 1

O comando a seguir iniciará um produtor para que você possa gravar três registros no tópico prodlistener.

kafka-console-producer \
    --broker-list localhost:9092 \
    --topic prodlistener

No console, você pode inserir esses três registros no tópico. Depois que isso for feito, esses registros serão confirmados no contêiner do Azure Cosmos DB for NoSQL mapeado para o tópico (produtos).

{"id": "0ac8b014-c3f4-4db0-8a1f-434bab460938", "name": "handlebar", "categoryId": "78148556-4e84-44be-abae-9755dde9c9e3"}
{"id": "54ba00da-50cf-44d8-b122-1d18bd1db400", "name": "handlebar", "categoryId": "eb642a5e-0c6f-4c83-b96b-bb2903b85e59"}
{"id": "381dde84-e6c2-4583-b66c-e4a4116f7d6e", "name": "handlebar", "categoryId": "cf8ae707-6d74-4563-831a-06e15a70a0dc"}

Ler do Azure Cosmos DB

Você pode criar um conector de origem no Kafka Connect usando um objeto de configuração JSON. Nesta configuração de exemplo abaixo, a maioria das propriedades deve permanecer inalterada, mas não se esqueça de alterar os seguintes valores:

Propriedade Descrição
connect.cosmos.connection.endpoint URI do ponto de extremidade da conta real
connect.cosmos.master.key Sua chave de conta real
connect.cosmos.databasename O nome do recurso de banco de dados da conta real
connect.cosmos.containers.topicmap Usando o formato CSV, um mapeamento de seus tópicos reais do Kafka para contêineres
{
  "name": "cosmosdb-source-connector",
  "config": {
    "connector.class": "com.azure.cosmos.kafka.connect.source.CosmosDBSourceConnector",
    "tasks.max": "1",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "connect.cosmos.task.poll.interval": "100",
    "connect.cosmos.connection.endpoint": "<cosmos-endpoint>",
    "connect.cosmos.master.key": "<cosmos-key>",
    "connect.cosmos.databasename": "<cosmos-database>",
    "connect.cosmos.containers.topicmap": "<kafka-topic>#<cosmos-container>",
    "connect.cosmos.offset.useLatest": false,
    "value.converter.schemas.enable": "false",
    "key.converter.schemas.enable": "false"
  }
}

Como ilustração, usando esta tabela de configuração de exemplo:

Propriedade Descrição
connect.cosmos.connection.endpoint https://dp420.documents.azure.com:443/
connect.cosmos.master.key C2y6yDjf5/R+ob0N8A7Cgv30VRDJIWEHLM+4QDU5DE2nQ9nDuVTqobD4b8mGGyPMbIZnqyMsEcaGQy67XIw/Jw==
connect.cosmos.databasename cosmicworks
connect.cosmos.containers.topicmap prodlistener#products

Aqui está um arquivo de configuração de exemplo:

{
  "name": "cosmosdb-source-connector",
  "config": {
    "connector.class": "com.azure.cosmos.kafka.connect.source.CosmosDBSourceConnector",
    "tasks.max": "1",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "connect.cosmos.task.poll.interval": "100",
    "connect.cosmos.connection.endpoint": "https://dp420.documents.azure.com:443/",
    "connect.cosmos.master.key": "C2y6yDjf5/R+ob0N8A7Cgv30VRDJIWEHLM+4QDU5DE2nQ9nDuVTqobD4b8mGGyPMbIZnqyMsEcaGQy67XIw/Jw==",
    "connect.cosmos.databasename": "cosmicworks",
    "connect.cosmos.containers.topicmap": "prodlistener#products",
    "connect.cosmos.offset.useLatest": false,
    "value.converter.schemas.enable": "false",
    "key.converter.schemas.enable": "false"
  }
}

Depois de configurados, o feed de alterações do Azure Cosmos DB será publicado em um tópico do Kafka.