Mover dados usando um conector Kafka
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.