Spostare i dati usando il connettore Spark per Azure Cosmos DB
Con Azure Synapse Analytics e Collegamento ad Azure Synapse per Azure Cosmos DB, è possibile creare un'elaborazione transazionale e analitica ibrida nativa del cloud per eseguire analisi sui dati in Azure Cosmos DB per NoSQL. Questa connessione consente l'integrazione sulla pipeline di dati tra le due estremità dell'ambiente dati, Azure Cosmos DB e Azure Synapse Analytics.
Attrezzaggio
Prima di tutto, assicurarsi che Collegamento a Synapse sia abilitato a livello di account. Questa operazione può essere eseguita tramite il portale di Azure o l'interfaccia della riga di comando di Azure:
az cosmosdb create --name <name> --resource-group <resource-group> --enable-analytical-storage true
È anche possibile usare Azure PowerShell:
New-AzCosmosDBAccount -ResourceGroupName <resource-group> -Name <name> -Location <location> -EnableAnalyticalStorage true
Quando si crea un contenitore, è necessario abilitare l'archivio analitico a livello di contenitore e per contenitore. Di nuovo, questa operazione si può eseguire con il portale.
Si può anche utilizzare l'interfaccia della riga di comando:
az cosmosdb sql container create --resource-group <resource-group> --account <account> --database <database> --name <name> --partition-key-path <partition-key-path> --throughput <throughput> --analytical-storage-ttl -1
Oppure Azure PowerShell:
New-AzCosmosDBSqlContainer -ResourceGroupName <resource-group> -AccountName <account> -DatabaseName <database> -Name <name> -PartitionKeyPath <partition-key-path> -Throughput <throughput> -AnalyticalStorageTtl -1
Suggerimento
È anche possibile usare i vari SDK per sviluppatori per abilitare o disabilitare l'archivio analitico a livello di contenitore o Collegamento a Synapse a livello di account.
Leggere in Azure Cosmos DB
Nota
I prossimi due esempi di Python devono essere eseguiti all'interno dell'area di lavoro di Azure Synapse Analytics.
Esistono due opzioni per eseguire query sui dati da Azure Cosmos DB for NoSQL. La prima è scegliere di caricare in un dataframe Spark in cui i metadati vengono memorizzati nella cache. In questo esempio viene usato Python per caricare un dataframe Spark che punta a un account Azure Cosmos DB for NoSQL.
productsDataFrame = spark.read.format("cosmos.olap")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.load()
In alternativa, è possibile creare una tabella Spark che punta direttamente ad Azure Cosmos DB for NoSQL. È quindi possibile eseguire query SparkSQL sulla tabella Spark senza incidere sull'archivio sottostante. In questo esempio viene usato Python per creare una tabella Spark.
create table products_qry using cosmos.olap options (
spark.synapse.linkedService 'cosmicworks_serv',
spark.cosmos.container 'products'
)
Scrivere in Azure Cosmos DB
Nota
I prossimi due esempi di Python devono essere eseguiti all'interno dell'area di lavoro di Azure Synapse Analytics.
Se si vogliono scrivere nuovi dati in Azure Cosmos DB dal dataframe Spark, è possibile usare lo script Python seguente per accodare i dati in un dataframe a un contenitore esistente.
productsDataFrame.write.format("cosmos.oltp")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.mode('append')\
.save()
Nota
Questa operazione inciderà sui carichi di lavoro delle transazioni esistenti e utilizzerà unità richiesta nei contenitori di Azure Cosmos DB for NoSQL.
Si può anche andare oltre e trasmettere dati da un dataframe, a partire da un checkpoint. È anche possibile accodare questi dati in streaming a un contenitore esistente usando questo script Python di esempio.
query = productsDataFrame\
.writeStream\
.format("cosmos.oltp")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.option("checkpointLocation", "/tmp/runIdentifier/")\
.outputMode("append")\
.start()
query.awaitTermination()