Przenoszenie danych przy użyciu łącznika Spark usługi Azure Cosmos DB

Ukończone

Za pomocą usług Azure Synapse Analytics i Azure Synapse Link dla usługi Azure Cosmos DB można utworzyć natywne dla chmury hybrydowe przetwarzanie transakcyjne i analityczne (HTAP), aby uruchomić analizę danych w usłudze Azure Cosmos DB for NoSQL. To połączenie umożliwia integrację potoku danych na obu końcach świata danych, usług Azure Cosmos DB i Azure Synapse Analytics.

Ustawienia

Najpierw upewnij się, że usługa Synapse Link jest włączona na poziomie konta. Można to zrobić przy użyciu witryny Azure Portal lub przy użyciu interfejsu wiersza polecenia platformy Azure:

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

Możesz również użyć programu Azure PowerShell:

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

Podczas tworzenia kontenera należy włączyć magazyn analityczny na poziomie kontenera dla poszczególnych kontenerów. Można to zrobić ponownie za pomocą portalu.

Można to również wykonać za pomocą interfejsu wiersza polecenia:

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

Lub za pomocą programu Azure PowerShell:

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

Napiwek

Możesz również użyć różnych zestawów SDK deweloperów, aby włączyć lub wyłączyć magazyn analityczny na poziomie kontenera lub synapse Link na poziomie konta.

Odczyt z usługi Azure Cosmos DB

Uwaga

W obszarze roboczym usługi Azure Synapse Analytics należy wykonać kilka następnych przykładów języka Python.

Istnieją dwie opcje wykonywania zapytań dotyczących danych z usługi Azure Cosmos DB for NoSQL. Najpierw możesz załadować element do ramki danych Platformy Spark, w której metadane są buforowane. W tym przykładzie użyto języka Python do załadowania ramki danych Platformy Spark wskazującej konto usługi Azure Cosmos DB for NoSQL.

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

Alternatywnie możesz utworzyć tabelę Platformy Spark, która bezpośrednio wskazuje usługę Azure Cosmos DB for NoSQL. Następnie można uruchamiać zapytania SparkSQL względem tabeli Spark bez wpływu na bazowy magazyn. W tym przykładzie użyto języka Python do utworzenia tabeli Platformy Spark.

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

Zapisywanie w usłudze Azure Cosmos DB

Uwaga

W obszarze roboczym usługi Azure Synapse Analytics należy wykonać kilka następnych przykładów języka Python.

Jeśli chcemy zapisać nowe dane w usłudze Azure Cosmos DB z ramki danych Platformy Spark, możemy użyć następującego skryptu języka Python, aby dołączyć dane w ramce danych do istniejącego kontenera.

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

Uwaga

Ta operacja wpłynie na istniejące obciążenia transakcji i będzie zużywać jednostki żądań w kontenerze Usługi Azure Cosmos DB for NoSQL[s].

Możemy nawet pobrać je dalej i przesyłać strumieniowo dane z ramki danych, zaczynając od punktu kontrolnego. Możemy również dołączyć te dane przesyłane strumieniowo do istniejącego kontenera przy użyciu tego przykładowego skryptu języka 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()