Spostare i dati usando un connettore Kafka

Completato

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.