使用 Apache Spark 数据帧读取 OpenSharing 共享表

本文提供了使用 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>")
)