本文提供了使用 Apache Spark 查询使用 OpenSharing 共享的数据的语法示例。 使用 deltasharing 关键字作为数据帧操作的格式选项。
用于查询共享数据的其他选项
还可以在元存储中注册的 OpenSharing 目录中创建使用共享表名称的查询,例如以下示例中的查询:
SQL
SELECT * FROM shared_table_name
Python
spark.read.table("shared_table_name")
有关在 Azure Databricks 中配置 OpenSharing 以及使用共享表名称查询数据的详细信息,请参阅 读取使用 Databricks-to-Databricks OpenSharing 共享的数据(适用于接收方)。
可以使用结构化流式处理以增量方式处理共享表中的记录。 若要使用结构化流式处理,必须为表启用历史记录共享。 请参阅ALTER SHARE。 历史记录共享需要 Databricks Runtime 12.2 LTS 或更高版本。
如果共享表所在的源 Delta 表已启用更改数据馈送,并且对此共享启用了历史记录,则在使用结构化流式处理或批处理读取 OpenSharing 共享时,即可使用更改数据馈送。 请参阅在 Azure Databricks 中使用更改数据馈送。
使用 OpenSharing 格式关键字进行读取
Apache Spark 数据帧读取操作支持 deltasharing 关键字,如以下示例所示:
df = (spark.read
.format("deltasharing")
.load("<profile-path>#<share-name>.<schema-name>.<table-name>")
)
读取 OpenSharing 共享表的变更数据馈送
对于启用了历史记录共享和变更数据馈送的表,可以使用 Apache Spark 数据帧读取变更数据馈送记录。 历史记录共享需要 Databricks Runtime 12.2 LTS 或更高版本。
df = (spark.read
.format("deltasharing")
.option("readChangeFeed", "true")
.option("startingTimestamp", "2021-04-21 05:45:46")
.option("endingTimestamp", "2021-05-21 12:00:00")
.load("<profile-path>#<share-name>.<schema-name>.<table-name>")
)
使用结构化流处理读取 OpenSharing 共享表
对于具有共享历史的表,您可以使用共享表作为结构化流媒体的来源。 历史记录共享需要 Databricks Runtime 12.2 LTS 或更高版本。
streaming_df = (spark.readStream
.format("deltasharing")
.load("<profile-path>#<share-name>.<schema-name>.<table-name>")
)
# If CDF is enabled on the source table
streaming_cdf_df = (spark.readStream
.format("deltasharing")
.option("readChangeFeed", "true")
.option("startingTimestamp", "2021-04-21 05:45:46")
.load("<profile-path>#<share-name>.<schema-name>.<table-name>")
)