Déplacer des données à l’aide d’un connecteur Kafka

Effectué

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.