Déplacer des données à l’aide d’un connecteur Kafka
Apache Kafka est une plateforme open source utilisée pour diffuser en continu des événements de manière distribuée. De nombreuses entreprises utilisent Kafka pour des scénarios d’intégration de données à hautes performances et à grande échelle. Kafka Connect est un outil dans sa suite permettant de diffuser en continu des données entre Kafka et d’autres systèmes de données. Bien entendu, cela peut inclure Azure Cosmos DB comme source de données ou cible (récepteur) de données.
Configuration
Les connecteurs Kafka Connect pour Azure Cosmos DB sont disponibles sous la forme d’un projet open source sur GitHub, à l’adresse microsoft/kafka-connect-cosmosdb. Les instructions de téléchargement et d’installation manuels du fichier JAR sont disponibles dans le référentiel.
Paramétrage
Quatre propriétés de configuration doivent être définies pour configurer correctement la connectivité à un compte Azure Cosmos DB for NoSQL.
| Propriété | Valeur |
|---|---|
| connect.cosmos.connection.endpoint | URI du point de terminaison de compte |
| connect.cosmos.master.key | Clé de compte |
| connect.cosmos.databasename | Nom de la ressource de base de données |
| connect.cosmos.containers.topicmap | Via le format CSV, mappage des rubriques Kafka aux conteneurs |
Mappage des rubriques aux conteneurs
Chaque conteneur doit être mappé à une rubrique. Par exemple, supposons que vous souhaitiez que le conteneur products (des produits) soit mappé à la rubrique prodlistener et que le conteneur customers (des clients) soit mappé à la rubrique custlistener. Dans ce cas, vous devez utiliser la chaîne de mappage CSV suivante : prodlistener#products,custlistener#customers.
Écrire dans Azure Cosmos DB
Nous allons écrire des données dans Azure Cosmos DB en créant une rubrique. Dans Apache Kafka, tous les messages sont envoyés via des rubriques.
Vous pouvez créer une nouvelle rubrique à l’aide de la commande kafka-topics. Cet exemple crée une nouvelle rubrique nommée prodlistener.
kafka-topics --create \
--zookeeper localhost:2181 \
--topic prodlistener \
--replication-factor 1 \
--partitions 1
La commande suivante démarre un producteur vous permettant d’écrire trois enregistrements dans la rubrique prodlistener.
kafka-console-producer \
--broker-list localhost:9092 \
--topic prodlistener
Dans la console, vous pouvez ensuite entrer ces trois enregistrements dans la rubrique. Une fois cette opération effectuée, ces enregistrements sont validés dans le conteneur Azure Cosmos DB for NoSQL mappé à la rubrique (produits).
{"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"}
Lire à partir d’Azure Cosmos DB
Vous pouvez créer un connecteur source dans Kafka Connect à l’aide d’un objet de configuration JSON. Dans cet exemple de configuration ci-dessous, la plupart des propriétés doivent rester inchangées, mais veillez à modifier les valeurs suivantes :
| Propriété | Description |
|---|---|
| connect.cosmos.connection.endpoint | URI effectif de votre point de terminaison de compte |
| connect.cosmos.master.key | Clé effective de votre compte |
| connect.cosmos.databasename | Nom de votre ressource de base de données de compte effective |
| connect.cosmos.containers.topicmap | Via le format CSV, mappage de vos rubriques Kafka effectives aux conteneurs |
{
"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"
}
}
En guise d’exemple, voyons l’utilisation de cette table de configuration :
| Propriété | Description |
|---|---|
| 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 |
Voici un exemple de fichier de configuration :
{
"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"
}
}
Une fois ce fichier configuré, les données du flux de modification Azure Cosmos DB sont publiées dans une rubrique Kafka.