Gegevens verplaatsen met behulp van een Kafka-connector
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.