Flytta data med hjälp av en Kafka-anslutning

Slutförd

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.