使用 Azure Cosmos DB Spark 連接器移動資料

已完成

使用適用於 Azure Cosmos DB 的 Azure Synapse Analytics 和 Azure Synapse Link,您可以建立雲端原生的混合式交易和分析處理 (HTAP),以在適用於 NoSQL 的 Azure Cosmos DB 中對您的數據執行分析。 如此連接後,即可整合您兩邊資料環境 (Azure Cosmos DB 和 Azure Synapse Analytics) 的資料管線。

設定

首先,您應該確定已在帳戶層級啟用 Synapse Link 。 可以使用 Azure 入口網站或使用 Azure CLI 完成此動作:

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

您也可以使用 Azure PowerShell:

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

建立容器時,您應以容器為單位,在容器層級啟用分析儲存體。 此動作同樣可以使用入口網站完成。

這也可以使用 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

或者,使用 Azure PowerShell 來完成:

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

提示

您也可以使用各種開發人員 SDK,在容器層級啟用或停用分析儲存體,或在帳戶層級啟用或停用 Synapse Link。

從 Azure Cosmos DB 讀取

注意

接下來的幾個 Python 範例應在您的 Azure Synapse Analytics 工作區內執行。

要從 Azure Cosmos DB for NoSQL 查詢資料時,有兩個選項。 首先,您可以選擇載入到中繼資料快取所在的 Spark DataFrame 中。 此範例會使用 Python,載入指向 Azure Cosmos DB for NoSQL 帳戶的 Spark DataFrame。

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

或者,您可以建立直接指向 Azure Cosmos DB for NoSQL 的 Spark 資料表。 接著,您可以對 Spark 資料表執行 SparkSQL 查詢,而不會影響底層存放區。 此範例會使用 Python 來建立 Spark 資料表。

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

寫入至 Azure Cosmos DB

注意

接下來的幾個 Python 範例應在您的 Azure Synapse Analytics 工作區內執行。

如果我們想要從 Spark DataFrame 將新資料寫入 Azure Cosmos DB,可以使用下列 Python 指令碼,將 DataFrame 中的資料附加至現有的容器。

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

注意

此作業會影響我們現有的交易工作負載,並取用 Azure Cosmos DB for NoSQL 容器上的要求單位。

我們甚至可以進一步從檢查點開始,自 DataFrame 串流資料。 我們也可以使用此範例 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()