Verschieben von Daten mithilfe eines Kafka-Connectors
Apache Kafka ist eine Open-Source-Plattform, die verwendet wird, um Ereignisse auf verteilte Weise zu streamen. Viele Unternehmen verwenden Kafka für umfangreiche Hochleistungs-Datenintegrationsszenarien. Kafka Connect ist ein Tool in ihrer Suite, um Daten zwischen Kafka und anderen Datensystemen zu streamen. Verständlicherweise kann dies Azure Cosmos DB als Datenquelle oder Datenziel (Senke) umfassen.
Einrichtung
Die Kafka Connect-Connectors für Azure Cosmos DB sind als Open-Source-Projekt auf GitHub unter microsoft/kafka-connect-cosmosdb verfügbar. Anweisungen zum manuellen Herunterladen und Installieren der JAR-Datei finden Sie im Repository.
Konfiguration
Es müssen vier Konfigurationseigenschaften festgelegt werden, um Konnektivität mit einem Azure Cosmos DB for NoSQL-Konto ordnungsgemäß zu konfigurieren.
| Eigentum | Wert |
|---|---|
| connect.cosmos.connection.endpoint | Kontoendpunkt-URI |
| connect.cosmos.master.key | Kontoschlüssel |
| connect.cosmos.databasename | Name der Datenbankressource |
| connect.cosmos.containers.topicmap | Zuordnung der Kafka-Themen zu Containern im CSV-Format |
Zuordnung von Themen zu Containern
Jeder Container sollte einem Thema zugeordnet werden. Angenommen, Sie möchten, dass der Produktcontainer dem Produktlistenerthema und dem Kundencontainer zum Thema "Custlistener " zugeordnet werden soll. In diesem Fall sollten Sie die folgende CSV-Zuordnungszeichenfolge verwenden: prodlistener#products,custlistener#customers.
Schreiben in Azure Cosmos DB
Nun werden durch Erstellen eines Themas Daten in Azure Cosmos DB geschrieben. In Apache Kafka werden alle Nachrichten über Themen gesendet.
Sie können ein neues Thema mit dem Befehl "kafka-topics " erstellen. In diesem Beispiel wird ein neues Thema namens prodlistener erstellt.
kafka-topics --create \
--zookeeper localhost:2181 \
--topic prodlistener \
--replication-factor 1 \
--partitions 1
Mit dem folgenden Befehl wird ein Producer gestartet, sodass Sie drei Datensätze in das Thema „prodlistener“ schreiben können.
kafka-console-producer \
--broker-list localhost:9092 \
--topic prodlistener
In der Konsole können Sie dann diese drei Datensätze in das Thema eingeben. Sobald dies geschehen ist, werden diese Datensätze dem Azure Cosmos DB für NoSQL-Container zugeordnet, der dem Thema (Produkte) zugeordnet ist.
{"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"}
Lesen aus Azure Cosmos DB
Sie können einen Quellconnector in Kafka Connect mithilfe eines JSON-Konfigurationsobjekts erstellen. In der folgenden Beispielkonfiguration sollten die meisten Eigenschaften unverändert bleiben, die folgenden Werte jedoch unbedingt angegeben werden:
| Eigentum | Beschreibung |
|---|---|
| connect.cosmos.connection.endpoint | Ihr tatsächlicher Kontoendpunkt-URI |
| connect.cosmos.master.key | Ihr tatsächlicher Kontoschlüssel |
| connect.cosmos.databasename | Der Namen Ihrer tatsächlichen Kontodatenbankressource |
| connect.cosmos.containers.topicmap | Eine Zuordnung Ihrer tatsächlichen Kafka-Themen zu Containern im CSV-Format |
{
"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"
}
}
Zur Veranschaulichung wird die folgende Beispielkonfigurationstabelle verwendet:
| Eigentum | Beschreibung |
|---|---|
| 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 |
Hier sehen Sie eine Beispielkonfigurationsdatei:
{
"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"
}
}
Nach der Konfiguration werden Daten aus dem Azure Cosmos DB-Änderungsfeed in einem Kafka-Thema veröffentlicht.