Verschieben von Daten mithilfe eines Kafka-Connectors

Abgeschlossen

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.