Mover dados usando o conector Spark do Azure Cosmos DB

Concluído

Com o Azure Synapse Analytics e o Link do Azure Synapse para Azure Cosmos DB, você pode criar um HTAP (processamento transacional e analítico) híbrido nativo de nuvem 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 seu mundo de dados, Azure Cosmos DB e Azure Synapse Analytics.

Instalação

Primeiro, verifique se o Link do Synapse está habilitado 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. Novamente, 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

Dica

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 Link do Synapse no nível da conta.

Ler do Azure Cosmos DB

Observação

Os próximos exemplos do Python devem ser executados em seu workspace do Azure Synapse Analytics.

Há duas opções para consultar dados no Azure Cosmos DB for NoSQL. Primeiro, você pode optar por carregar em um DataFrame do Spark em que os metadados são armazenados em cache. Este exemplo usa o Python para carregar um DataFrame do Spark que aponta para uma conta no 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, você pode criar uma tabela do Spark que aponta ao Azure Cosmos DB for NoSQL diretamente. Em seguida, você pode executar consultas do SparkSQL na tabela Spark sem afetar o armazenamento subjacente. Este exemplo usa o Python para criar uma tabela do Spark.

create table products_qry using cosmos.olap options (
    spark.synapse.linkedService 'cosmicworks_serv',
    spark.cosmos.container 'products'
)

Gravar no Azure Cosmos DB

Observação

Os próximos exemplos do Python devem ser executados em seu workspace do Azure Synapse Analytics.

Para gravar novos dados no Azure Cosmos DB do nosso DataFrame do Spark, podemos usar o script Python a seguir para acrescentar os dados em um DataFrame a um contêiner.

productsDataFrame.write.format("cosmos.oltp")\
    .option("spark.synapse.linkedService", "cosmicworks_serv")\
    .option("spark.cosmos.container", "products")\
    .mode('append')\
    .save()

Observação

Essa operação afetará as cargas de trabalho de transação existentes e consumirá unidades de solicitação nos contêineres no Azure Cosmos DB for NoSQL.

Podemos até mesmo ir além e transmitir dados de um DataFrame, começando em um ponto de verificação. Também podemos acrescentar esses dados de streaming a um contêiner 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()