Flytta data med hjälp av en Kafka-anslutning
Apache Kafka är en plattform med öppen källkod som används för att strömma händelser på ett distribuerat sätt. Många företag använder Kafka för storskaliga scenarier för dataintegrering med höga prestanda. Kafka Connect är ett verktyg i deras svit för att strömma data mellan Kafka och andra datasystem. Det är förståeligt att detta kan inkludera Azure Cosmos DB som en datakälla eller ett mål (mottagare) av data.
Ställ in
Kafka Connect-anslutningsapparna för Azure Cosmos DB är tillgängliga som ett projekt med öppen källkod på GitHub på microsoft/kafka-connect-cosmosdb. Instruktioner för att ladda ned och installera JAR-filen manuellt finns på lagringsplatsen.
Konfiguration
Fyra konfigurationsegenskaper bör anges för att konfigurera anslutningen till ett Azure Cosmos DB för NoSQL-konto korrekt.
| Egenskap | Värde |
|---|---|
| connect.cosmos.connection.endpoint | Kontoslutpunkts-URI |
| connect.cosmos.master.key | Kontonyckel |
| connect.cosmos.databasename | Namnet på databasresursen |
| connect.cosmos.containers.topicmap | Med HJÄLP av CSV-format mappas Kafka-ämnen till containrar |
Mappa ämnen till containrar
Varje container ska mappas till ett ämne. Anta till exempel att du vill att produktcontainern ska mappas till ämnet prodlistener och kundcontainern till ämnet custlistener . I så fall bör du använda följande CSV-mappningssträng: prodlistener#products,custlistener#customers.
Skriva till Azure Cosmos DB
Nu ska vi skriva data till Azure Cosmos DB genom att skapa ett ämne. I Apache Kafka skickas alla meddelanden via ämnen.
Du kan skapa ett nytt ämne med hjälp av kommandot kafka-topics . I det här exemplet skapas ett nytt ämne med namnet prodlistener.
kafka-topics --create \
--zookeeper localhost:2181 \
--topic prodlistener \
--replication-factor 1 \
--partitions 1
Följande kommando startar en producent så att du kan skriva tre poster till prodlistener-ämnet.
kafka-console-producer \
--broker-list localhost:9092 \
--topic prodlistener
Och i -konsolen kan du sedan ange dessa tre poster i ämnet. När detta är klart kommer dessa poster att registreras i Azure Cosmos DB för NoSQL-containern som är kopplad till ämneskategorin (produkter).
{"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"}
Läsa från Azure Cosmos DB
Du kan skapa en källanslutning i Kafka Connect med hjälp av ett JSON-konfigurationsobjekt. I den här exempelkonfigurationen nedan bör de flesta egenskaperna lämnas oförändrade, men se till att ändra följande värden:
| Egenskap | Beskrivning |
|---|---|
| connect.cosmos.connection.endpoint | Din faktiska kontoslutpunkts-URI |
| connect.cosmos.master.key | Din faktiska kontonyckel |
| connect.cosmos.databasename | Namnet på din faktiska kontodatabasresurs |
| connect.cosmos.containers.topicmap | Med CSV-format kan du mappa dina faktiska Kafka-ämnen till containrar |
{
"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"
}
}
Som ett illustrativt exempel använder du den här exempelkonfigurationstabellen:
| Egenskap | Beskrivning |
|---|---|
| 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 |
Här är ett exempel på en konfigurationsfil:
{
"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"
}
}
När de har konfigurerats publiceras data från Azure Cosmos DB-ändringsflödet till ett Kafka-ämne.