Gegevens verplaatsen met behulp van een Kafka-connector

Voltooid

Apache Kafka is een opensource-platform dat wordt gebruikt om gebeurtenissen op een gedistribueerde manier te streamen. Veel bedrijven gebruiken Kafka voor grootschalige high-performance gegevensintegratiescenario's. Kafka Connect is een hulpprogramma in hun suite om gegevens te streamen tussen Kafka en andere gegevenssystemen. Begrijpelijk, dit kan Azure Cosmos DB omvatten als gegevensbron of een doel (sink) met gegevens.

Instellingen

De Kafka Connect-connectors voor Azure Cosmos DB zijn beschikbaar als een opensource-project op GitHub op microsoft/kafka-connect-cosmosdb. Instructies voor het handmatig downloaden en installeren van het JAR-bestand zijn beschikbaar in de opslagplaats.

Configuratie

Er moeten vier configuratie-eigenschappen worden ingesteld om de connectiviteit met een Azure Cosmos DB for NoSQL-account correct te configureren.

Eigenschap Waarde
connect.cosmos.connection.endpoint Eindpunt-URI van account
connect.cosmos.master.key Accountsleutel
connect.cosmos.databasename Naam van de databaseresource
connect.cosmos.containers.topicmap Csv-indeling gebruiken, een toewijzing van de Kafka-onderwerpen aan containers

Onderwerpen voor containers toewijzen

Elke container moet worden toegewezen aan een onderwerp. Stel dat u wilt dat de container producten wordt toegewezen aan het onderwerp prodlistener en de container klanten aan het onderwerp custlistener . In dat geval moet u de volgende CSV-toewijzingstekenreeks gebruiken: prodlistener#products,custlistener#customers.

Schrijven naar Azure Cosmos DB

Laten we gegevens schrijven naar Azure Cosmos DB door een onderwerp te maken. In Apache Kafka worden alle berichten verzonden via onderwerpen.

U kunt een nieuw onderwerp maken met behulp van de opdracht kafka-topics . In dit voorbeeld wordt een nieuw onderwerp met de naam prodlistener gemaakt.

kafka-topics --create \
    --zookeeper localhost:2181 \
    --topic prodlistener \
    --replication-factor 1 \
    --partitions 1

Met de volgende opdracht wordt een producent gestart, zodat u drie records naar het onderwerp prodlistener kunt schrijven.

kafka-console-producer \
    --broker-list localhost:9092 \
    --topic prodlistener

En in de console kunt u deze drie records invoeren in het onderwerp. Zodra dit is gebeurd, worden deze records doorgevoerd in de Azure Cosmos DB for NoSQL-container die is toegewezen aan het onderwerp (producten).

{"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"}

Lezen uit Azure Cosmos DB

U kunt een bronconnector maken in Kafka Connect met behulp van een JSON-configuratieobject. In deze voorbeeldconfiguratie hieronder moeten de meeste eigenschappen ongewijzigd blijven, maar zorg ervoor dat u de volgende waarden wijzigt:

Eigenschap Beschrijving
connect.cosmos.connection.endpoint De werkelijke eindpunt-URI van uw account
connect.cosmos.master.key Uw werkelijke accountsleutel
connect.cosmos.databasename De naam van uw werkelijke accountdatabaseresource
connect.cosmos.containers.topicmap Csv-indeling gebruiken, een toewijzing van uw werkelijke Kafka-onderwerpen aan containers
{
  "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"
  }
}

Als illustratief voorbeeld gebruikt u deze voorbeeldconfiguratietabel:

Eigenschap Beschrijving
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 volgt een voorbeeld van een configuratiebestand:

{
  "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"
  }
}

Zodra deze is geconfigureerd, worden gegevens uit de Wijzigingenfeed van Azure Cosmos DB gepubliceerd naar een Kafka-onderwerp.