Movimiento de datos mediante el conector spark de Azure Cosmos DB
Con Azure Synapse Analytics y Azure Synapse Link para Azure Cosmos DB, puede crear un procesamiento transaccional y analítico híbrido nativo en la nube (HTAP) para ejecutar análisis sobre los datos de Azure Cosmos DB para NoSQL. Esta conexión permite la integración mediante la canalización de datos en ambos extremos del mundo de los datos, Azure Cosmos DB y Azure Synapse Analytics.
Configurar
En primer lugar, debe asegurarse de que Synapse Link está habilitado en el nivel de cuenta. Esto se puede lograr mediante el Azure Portal o mediante el CLI de Azure:
az cosmosdb create --name <name> --resource-group <resource-group> --enable-analytical-storage true
También puede usar Azure PowerShell:
New-AzCosmosDBAccount -ResourceGroupName <resource-group> -Name <name> -Location <location> -EnableAnalyticalStorage true
Al crear un contenedor, debe habilitar el almacenamiento analítico en el nivel de contenedor de uno en uno. De nuevo, esto se puede lograr con el portal.
Esto también se puede lograr con la 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
O con Azure PowerShell:
New-AzCosmosDBSqlContainer -ResourceGroupName <resource-group> -AccountName <account> -DatabaseName <database> -Name <name> -PartitionKeyPath <partition-key-path> -Throughput <throughput> -AnalyticalStorageTtl -1
Sugerencia
También puede usar los distintos SDK de desarrollador para habilitar o deshabilitar el almacenamiento analítico a nivel de contenedor o Synapse Link en el nivel de cuenta.
Lectura desde Azure Cosmos DB
Nota:
El siguiente par de ejemplos de Python debe realizarse en el área de trabajo de Azure Synapse Analytics.
Hay dos opciones para consultar datos desde Azure Cosmos DB for NoSQL. En primer lugar, puede elegir cargar un DataFrame de Spark donde se almacenan en caché los metadatos. En este ejemplo, se usa Python para cargar un DataFrame de Spark que apunta a una cuenta de Azure Cosmos DB for NoSQL.
productsDataFrame = spark.read.format("cosmos.olap")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.load()
Como alternativa, puede crear una tabla de Spark que apunte directamente a la instancia de Azure Cosmos DB for NoSQL. A continuación, puede ejecutar consultas de SparkSQL en la tabla de Spark sin afectar al almacén subyacente. En este ejemplo se usa Python para crear una tabla de Spark.
create table products_qry using cosmos.olap options (
spark.synapse.linkedService 'cosmicworks_serv',
spark.cosmos.container 'products'
)
Escritura a Azure Cosmos DB
Nota:
El siguiente par de ejemplos de Python debe realizarse en el área de trabajo de Azure Synapse Analytics.
Si queremos escribir datos nuevos en Azure Cosmos DB desde nuestro DataFrame de Spark, podemos usar el siguiente script de Python para anexar los datos de un DataFrame a un contenedor existente.
productsDataFrame.write.format("cosmos.oltp")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.mode('append')\
.save()
Nota:
Esta operación afectará a las cargas de trabajo de transacciones existentes y consumirá unidades de solicitud en los contenedores de Azure Cosmos DB for NoSQL.
Incluso podemos ir más allá y transmitir datos desde un DataFrame, empezando desde un punto de control. También podemos anexar estos datos de streaming a un contenedor existente mediante este script de Python de ejemplo.
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()