Flytta data med hjälp av Azure Cosmos DB Spark-anslutningsappen

Slutförd

Med Azure Synapse Analytics och Azure Synapse Link för Azure Cosmos DB kan du skapa en molnbaserad hybridtransaktions- och analysbearbetning (HTAP) för att köra analys över dina data i Azure Cosmos DB för NoSQL. Den här anslutningen möjliggör integrering över din datapipeline i båda ändar av din datavärld, Azure Cosmos DB och Azure Synapse Analytics.

Ställ in

Först bör du se till att Synapse Link är aktiverat på kontonivå. Detta kan göras med hjälp av Azure Portal eller med hjälp av Azure CLI:

az cosmosdb create --name <name> --resource-group <resource-group> --enable-analytical-storage true

Du kan också använda Azure PowerShell:

New-AzCosmosDBAccount -ResourceGroupName <resource-group> -Name <name>  -Location <location> -EnableAnalyticalStorage true

När du skapar en container bör du aktivera analyslagring på containernivå per container. Återigen kan detta åstadkommas med portalen.

Detta kan också göras med 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

Eller med Azure PowerShell:

New-AzCosmosDBSqlContainer -ResourceGroupName <resource-group> -AccountName <account> -DatabaseName <database> -Name <name> -PartitionKeyPath <partition-key-path> -Throughput <throughput> -AnalyticalStorageTtl -1

Dricks

Du kan också använda de olika utvecklarnas SDK:er för att aktivera eller inaktivera antingen analytisk lagring på containernivå eller Synapse Link på kontonivå.

Läsa från Azure Cosmos DB

Kommentar

Nästa par Python-exempel bör utföras på din Azure Synapse Analytics-arbetsyta.

Det finns två alternativ för att fråga efter data från Azure Cosmos DB för NoSQL. Först kan du välja att läsa in till en Spark DataFrame där metadata cachelagras. I det här exemplet används Python för att läsa in en Spark DataFrame som pekar på ett Azure Cosmos DB för NoSQL-konto.

productsDataFrame = spark.read.format("cosmos.olap")\
    .option("spark.synapse.linkedService", "cosmicworks_serv")\
    .option("spark.cosmos.container", "products")\
    .load()

Du kan också skapa en Spark-tabell som pekar direkt på Azure Cosmos DB för NoSQL. Du kan sedan köra SparkSQL-frågor mot Spark-tabellen utan att påverka det underliggande arkivet. I det här exemplet används Python för att skapa en Spark-tabell.

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

Skriva till Azure Cosmos DB

Kommentar

Nästa par Python-exempel bör utföras på din Azure Synapse Analytics-arbetsyta.

Om vi vill skriva nya data till Azure Cosmos DB från vår Spark DataFrame kan vi använda följande Python-skript för att lägga till data i en DataFrame till en befintlig container.

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

Kommentar

Den här åtgärden påverkar våra befintliga transaktionsarbetsbelastningar och förbrukar enheter för begäranden i Azure Cosmos DB for NoSQL-containern[s].

Vi kan till och med ta det vidare och strömma data från en DataFrame, med början från en kontrollpunkt. Vi kan också lägga till dessa strömmande data i en befintlig container med hjälp av det här Python-exempelskriptet.

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()