Déplacer des données à l’aide du connecteur Spark Azure Cosmos DB
Avec Azure Synapse Analytics et Azure Synapse Link pour Azure Cosmos DB, vous pouvez créer un traitement transactionnel et analytique hybride natif cloud (HTAP) pour exécuter l’analytique sur vos données dans Azure Cosmos DB pour NoSQL. Cette connexion facilite l’intégration de votre pipeline de données aux deux extrémités de votre environnement de données, Azure Cosmos DB et Azure Synapse Analytics.
Programme d’installation
Tout d’abord, vous devez vous assurer que Synapse Link est activé au niveau du compte. Pour ce faire, vous pouvez utiliser le portail Azure ou l’interface Azure CLI :
az cosmosdb create --name <name> --resource-group <resource-group> --enable-analytical-storage true
Vous pouvez également utiliser Azure PowerShell :
New-AzCosmosDBAccount -ResourceGroupName <resource-group> -Name <name> -Location <location> -EnableAnalyticalStorage true
Lorsque vous créez un conteneur, vous devez activer le stockage analytique au niveau du conteneur pour chaque conteneur individuellement. Là encore, cette opération peut être effectuée via le portail.
Elle peut également être accomplie via l’interface CLI :
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
Ou avec Azure PowerShell :
New-AzCosmosDBSqlContainer -ResourceGroupName <resource-group> -AccountName <account> -DatabaseName <database> -Name <name> -PartitionKeyPath <partition-key-path> -Throughput <throughput> -AnalyticalStorageTtl -1
Conseil
Vous pouvez également utiliser les différents kits SDK pour développeurs afin d’activer ou de désactiver le stockage analytique au niveau de chaque conteneur ou Synapse Link au niveau du compte.
Lire à partir d’Azure Cosmos DB
Notes
Les deux exemples Python suivants doivent être exécutés dans votre espace de travail Azure Synapse Analytics.
Il existe deux options pour interroger des données à partir d’Azure Cosmos DB for NoSQL. Tout d’abord, vous pouvez choisir de charger un DataFrame Spark où les métadonnées sont mises en cache. Cet exemple utilise Python pour charger un DataFrame Spark qui pointe vers un compte Azure Cosmos DB for NoSQL.
productsDataFrame = spark.read.format("cosmos.olap")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.load()
Vous pouvez également créer une table Spark qui pointe directement vers Azure Cosmos DB for NoSQL. Vous pouvez ensuite exécuter des requêtes SparkSQL sur cette table Spark sans affecter le magasin sous-jacent. Cet exemple utilise Python pour créer une table Spark.
create table products_qry using cosmos.olap options (
spark.synapse.linkedService 'cosmicworks_serv',
spark.cosmos.container 'products'
)
Écrire dans Azure Cosmos DB
Notes
Les deux exemples Python suivants doivent être exécutés dans votre espace de travail Azure Synapse Analytics.
Si vous souhaitez écrire de nouvelles données dans Azure Cosmos DB à partir de votre DataFrame Spark, vous pouvez utiliser le script Python suivant pour ajouter les données d’un DataFrame dans un conteneur existant.
productsDataFrame.write.format("cosmos.oltp")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.mode('append')\
.save()
Notes
Cette opération a un impact sur nos charges de travail de transaction existantes et consomme des unités de requête sur le ou les conteneurs Azure Cosmos DB for NoSQL.
Vous pouvez même aller plus loin et diffuser en continu des données à partir d’un DataFrame, en commençant à partir d’un point de contrôle. Vous pouvez également ajouter ces données de streaming dans un conteneur existant en utilisant cet exemple de script Python.
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()