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