Mover dados usando um conector Kafka

Concluído

Apache Kafka é uma plataforma de código aberto usada para transmitir eventos de forma distribuída. Muitas empresas usam o Kafka para cenários de integração de dados de alto desempenho em grande escala. Kafka Connect é uma ferramenta dentro de seu pacote para transmitir dados entre Kafka e outros sistemas de dados. Compreensivelmente, isso pode incluir o Azure Cosmos DB como uma fonte de dados ou um destino (coletor) de dados.

Configurar

Os conectores Kafka Connect para Azure Cosmos DB estão disponíveis como um projeto de código aberto no GitHub em microsoft/kafka-connect-cosmosdb. 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 para NoSQL.

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

Tópicos para o mapa de 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 customers 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 através de tópicos.

Você pode criar um novo tópico usando o comando kafka-topics . Este exemplo criará um novo 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

E no console, você pode inserir esses três registros para o tópico. Feito isso, esses registros serão confirmados no contêiner do Azure Cosmos DB para 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 a partir 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 ser deixada inalterada, mas certifique-se de alterar os seguintes valores:

Propriedade Descrição
connect.cosmos.connection.endpoint URI do ponto de extremidade da sua conta real
connect.cosmos.master.key A 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 exemplo ilustrativo, 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 exemplo de arquivo de configuração:

{
  "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"
  }
}

Uma vez configurados, os dados do feed de alterações do Azure Cosmos DB serão publicados em um tópico do Kafka.