Spostare i dati usando un connettore Kafka
Apache Kafka è una piattaforma open source usata per trasmettere eventi in modo distribuito. Molte aziende usano Kafka per scenari di integrazione dei dati ad alte prestazioni su larga scala. Kafka Connect è uno strumento incluso nella suite che consente di trasmettere i dati tra Kafka e altri sistemi dati. I sistemi supportati includono chiaramente anche Azure Cosmos DB come origine dei dati o destinazione (sink) dei dati.
Attrezzaggio
I connettori di Kafka Connect per Azure Cosmos DB sono disponibili come progetto open source in GitHub in microsoft/kafka-connect-cosmosdb. Le istruzioni per il download e l'installazione manuale del file JAR sono disponibili nel repository.
Impostazione
Per configurare correttamente la connettività a un account di SQL di Azure Cosmos DB for NoSQL, sono necessarie quattro proprietà di configurazione.
| Proprietà | valore |
|---|---|
| connect.cosmos.connection.endpoint | URI dell'endpoint dell'account |
| connect.cosmos.master.key | Chiave dell'account |
| connect.cosmos.databasename | Nome della risorsa di database |
| connect.cosmos.containers.topicmap | Con formato CSV, mapping degli argomenti Kafka ai contenitori |
Argomenti per la mappa dei contenitori
Ogni contenitore deve essere mappato a un argomento. Si supponga, ad esempio, che il contenitore products sia mappato all'argomento prodlistener e al contenitore customers all'argomento custlistener. In questo caso, è consigliabile usare la stringa di mapping CSV seguente: prodlistener#products,custlistener#customers.
Scrivere in Azure Cosmos DB
Per scrivere i dati in Azure Cosmos DB è necessario creare un argomento. In Apache Kafka tutti i messaggi vengono inviati tramite argomenti.
È possibile creare un nuovo argomento usando il comando kafka-topics. In questo esempio verrà creato un nuovo argomento denominato prodlistener.
kafka-topics --create \
--zookeeper localhost:2181 \
--topic prodlistener \
--replication-factor 1 \
--partitions 1
Il comando seguente avvia un producer che consente di scrivere tre record nell'argomento prodlistener.
kafka-console-producer \
--broker-list localhost:9092 \
--topic prodlistener
Tramite la console è quindi possibile immettere questi tre record nell'argomento. Al termine, questi record verranno sottoposti a commit nel contenitore di SQL di Azure Cosmos DB for NoSQL mappato all'argomento (products).
{"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"}
Leggere in Azure Cosmos DB
È possibile creare un connettore di origine in Kafka Connect usando un oggetto di configurazione JSON. Nella configurazione di esempio che segue, la maggior parte delle proprietà deve essere lasciata invariata, ma assicurarsi di modificare i valori seguenti:
| Proprietà | Descrizione |
|---|---|
| connect.cosmos.connection.endpoint | URI dell'endpoint dell'account effettivo |
| connect.cosmos.master.key | Chiave dell'account effettiva |
| connect.cosmos.databasename | Nome della risorsa di database dell'account effettiva |
| connect.cosmos.containers.topicmap | Con formato CSV, mapping degli argomenti Kafka ai contenitori effettivi |
{
"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"
}
}
Come esempio dimostrativo, usando questa tabella di configurazione:
| Proprietà | Descrizione |
|---|---|
| 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 |
Di seguito è riportato un esempio di file di configurazione.
{
"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"
}
}
Dopo la configurazione, i dati dal feed di modifiche di Azure Cosmos DB verranno pubblicati in un argomento Kafka.