Mover dados usando o conector Azure Cosmos DB Spark
Com o Azure Synapse Analytics e o Azure Synapse Link para Azure Cosmos DB, você pode criar um processamento transacional e analítico híbrido nativo da nuvem (HTAP) para executar análises sobre seus dados no Azure Cosmos DB para NoSQL. Essa conexão permite a integração em seu pipeline de dados em ambas as extremidades do mundo de dados, Azure Cosmos DB e Azure Synapse Analytics.
Configurar
Primeiro, você deve certificar-se de que o Synapse Link está ativado no nível da conta. Isso pode ser feito usando o portal do Azure ou usando a CLI do Azure:
az cosmosdb create --name <name> --resource-group <resource-group> --enable-analytical-storage true
Você também pode usar o Azure PowerShell:
New-AzCosmosDBAccount -ResourceGroupName <resource-group> -Name <name> -Location <location> -EnableAnalyticalStorage true
Ao criar um contêiner, você deve habilitar o armazenamento analítico no nível do contêiner por contêiner. Mais uma vez, isso pode ser feito com o portal.
Isso também pode ser feito com a 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 com o Azure PowerShell:
New-AzCosmosDBSqlContainer -ResourceGroupName <resource-group> -AccountName <account> -DatabaseName <database> -Name <name> -PartitionKeyPath <partition-key-path> -Throughput <throughput> -AnalyticalStorageTtl -1
Gorjeta
Você também pode usar os vários SDKs de desenvolvedores para habilitar ou desabilitar o armazenamento analítico em um nível por contêiner ou Synapse Link no nível da conta.
Ler a partir do Azure Cosmos DB
Nota
Os próximos exemplos de Python devem ser executados em seu espaço de trabalho do Azure Synapse Analytics.
Há duas opções para consultar dados do Azure Cosmos DB para NoSQL. Primeiro, você pode optar por carregar em um Spark DataFrame onde os metadados são armazenados em cache. Este exemplo usa Python para carregar um Spark DataFrame que aponta para uma conta do Azure Cosmos DB para NoSQL.
productsDataFrame = spark.read.format("cosmos.olap")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.load()
Como alternativa, você pode criar uma tabela Spark que aponte diretamente para o Azure Cosmos DB para NoSQL. Em seguida, você pode executar consultas SparkSQL na tabela Spark sem afetar o armazenamento subjacente. Este exemplo usa Python para criar uma tabela Spark.
create table products_qry using cosmos.olap options (
spark.synapse.linkedService 'cosmicworks_serv',
spark.cosmos.container 'products'
)
Gravar no Azure Cosmos DB
Nota
Os próximos exemplos de Python devem ser executados em seu espaço de trabalho do Azure Synapse Analytics.
Se quisermos gravar novos dados no Azure Cosmos DB a partir do nosso Spark DataFrame, podemos usar o seguinte script Python para acrescentar os dados em um DataFrame a um contêiner existente.
productsDataFrame.write.format("cosmos.oltp")\
.option("spark.synapse.linkedService", "cosmicworks_serv")\
.option("spark.cosmos.container", "products")\
.mode('append')\
.save()
Nota
Essa operação afetará nossas cargas de trabalho de transação existentes e consumirá unidades de solicitação no contêiner do Azure Cosmos DB para NoSQL.
Podemos até levá-lo mais longe e transmitir dados de um DataFrame, começando a partir de um ponto de verificação. Também podemos acrescentar esses dados de streaming a um contêiner existente usando este script Python de exemplo.
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()