Spark 数据源

通过 Spark 数据源 API,可以直接从Azure Databricks读取和写入外部数据库。 仅当需要 Spark 引擎的完整灵活性、想要对源执行本机查询或需要对外部系统进行写入访问时,才使用它。 通常,Azure Databricks 建议采用受管控的只读访问,并支持自动将 Spark 或 SQL 查询下推。 请参阅什么是查询联邦?

Spark 数据源 API 具有连接性、查询执行和架构检测的特定行为。

  • 主工作负荷和任何后续 Spark 转换在 Azure Databricks Spark 群集上运行。
  • 使用 query 此选项时,指定的 SQL 语句完全在外部数据源上运行。 Spark 在不对查询字符串执行转换下推的情况下提取结果。
  • 连接需要Azure Databricks捆绑连接器、用户提供的 JDBC 驱动程序或 PySpark 自定义数据源。
  • Spark 会自动从外部数据库表读取架构,并将其类型映射到 Spark SQL 类型。

使用捆绑的连接器

Databricks Runtime 提供适用于常见数据源的经过优化的连接器。 有关完整列表,请参阅 支持的捆绑连接器

捆绑连接器使用 hostport 用作单独的选项,而不是完整的 JDBC URL 字符串。

使用直通查询读取数据

query使用此选项可确保在数据到达 Spark 之前对源数据库执行筛选和联接逻辑。 对于通过视图实现自动查询下推和 Unity Catalog 权限委派的受管控读取访问,请考虑改用 远程查询

df = (spark.read
  .format("sqlserver")
  .option("host", "<your-sql-server-instance>.database.chinacloudapi.cn")
  .option("user", dbutils.secrets.get(scope="<scope>", key="<user>"))
  .option("password", dbutils.secrets.get(scope="<scope>", key="<password>"))
  .option("database", "<database-name>")
  .option("query", "SELECT id, name FROM users WHERE active = 1")
  .load())

写入数据

使用 .mode() 指定写入模式,以控制数据的写入方式。 用于 append 向现有表添加行或 overwrite 替换其内容。

(df.write
  .format("sqlserver")
  .mode("overwrite")
  .option("host", "<your-sql-server-instance>.database.chinacloudapi.cn")
  .option("user", dbutils.secrets.get(scope="<scope>", key="<user>"))
  .option("password", dbutils.secrets.get(scope="<scope>", key="<password>"))
  .option("database", "<database-name>")
  .option("dbtable", "<table-name>")
  .save())

使用 JDBC UC 连接

如果源特定的连接器未捆绑,或者想要使用特定的 JDBC 驱动程序版本,请使用 JDBC Unity 目录连接。 这样,便可以集中凭据管理并自带 JDBC 驱动程序。

JDBC Unity 目录连接比直接使用捆绑连接器或原始 JDBC 驱动程序具有多种优势。 使用 JDBC Unity 目录连接,可以:

  • 为任何支持 JDBC 的数据库自备 JDBC 驱动程序 JAR 文件。
  • 创建连接一次,并在无服务器、标准和专用群集之间重复使用它。
  • 使用 Unity 目录连接对象利用对数据源的受管理访问权限。
  • 隐藏查询用户的连接凭据。
  • 通过 Spark 数据源 API 读取和写入外部数据库。

若要使用 JDBC Unity 目录连接,请在 Spark 选项中指定 databricks.connection

df = (spark.read
  .format("jdbc")
  .option("databricks.connection", "<connection-name>")
  .option("query", "SELECT * FROM external_table")
  .load())

有关设置说明,请参阅 JDBC 连接

在专用群集上使用自定义连接器

在专用(经典)群集上,可以安装未与 Databricks Runtime 捆绑的第三方 Spark 数据源连接器或 JDBC 驱动程序。

在以下情况下使用此方法:

  • 对于 MongoDB、Cassandra、Couchbase 或 Elasticsearch 等系统,需要第三方 Spark 连接器。
  • 你需要一个运行时中未包含的特定驱动程序版本。
  • 你希望直接在群集上安装 JDBC 驱动程序,而无需设置 Unity 目录连接。

安装连接器或驱动程序

通过计算>你的群集>>安装新库在群集上安装库。 可以直接使用 Maven 坐标,而无需下载或上传任何 JAR。 重启集群,以使该库生效。

读取数据

安装连接器后,请使用连接器的格式名称和所需的连接选项来读取数据。

df = (spark.read
  .format("mongodb")
  .option("connection.uri", "mongodb://<hostname>:27017")
  .option("database", "<database-name>")
  .option("collection", "<collection-name>")
  .load())

写入数据

使用相同的格式名称和连接选项将数据写回源。

(df.write
  .format("mongodb")
  .mode("overwrite")
  .option("connection.uri", "mongodb://<hostname>:27017")
  .option("database", "<database-name>")
  .option("collection", "<collection-name>")
  .save())

Considerations

在专用群集上使用自定义连接器时,请记住以下几点。

  • 驱动程序或连接器仅在安装驱动程序的群集上可用。
  • Databricks SQL、无服务器或标准访问模式群集不支持自定义第三方 Spark JAR。 对于这些计算类型,请使用捆绑连接器或 JDBC Unity 目录连接。

PySpark 自定义数据源

使用 Python DataSource API 可以完全在 Python 中生成自定义数据连接器,而无需使用 JAR 或基于 JVM 的库。 如果需要连接到没有 JDBC 接口的 REST API、SaaS 应用程序或任何系统,或者想要以编程方式生成综合数据,请使用此功能。 API 支持批处理和流式读取和写入。

注释

PySpark 自定义数据源需要 Databricks Runtime 15.4 LTS 或更高版本。

有关设置、示例和 API 参考,请参阅 PySpark 自定义数据源

比较集成策略

下表比较了 Spark 数据源 API、Lakehouse Federation 和 Lakeflow Connect,以帮助您根据自己的用例选择合适的方法。

功能 Spark 数据源 API 湖仓联合体 Lakeflow Connect
主要用例 复杂的 ETL、自定义 Spark 逻辑、直通查询 即席查询,BI 报告 大规模自动化引入
数据移动 加载到 Spark 内存(临时) 加载到 Spark 内存(临时) 复制到 Delta Lake (持久性)
查询执行 使用原生 query 选项手动下发 Spark 和 SQL 筛选器、联接和聚合的自动下推 不适用(整表复制)
治理 Unity 目录连接(JDBC)或机密范围 Unity Catalog(联邦目录) Unity Catalog(托管管道)
最适用于 需要 Spark 全部灵活性的高级用户 在保留治理的同时最大程度地减少数据移动 生产 CDC 和引入管道

支持的捆绑连接器

以下数据源捆绑在 Databricks Runtime 中,可以直接通过 Spark 调用。 专用群集和标准群集支持读取和写入。

注释

PostgreSQL、SQL Server、MySQL、Snowflake 和 Redshift 支持在无服务器计算上进行写入。 有关受支持的连接器选项,请参阅捆绑连接器的无服务器写入选项

数据源 spark.format() 名称
PostgreSQL "postgresql"
SQL Server "sqlserver"
MySQL 和 MariaDB "mysql"
Snowflake "snowflake"
Amazon Redshift "redshift"
谷歌BigQuery "bigquery"
Azure Synapse "SQLDW"
HTTP "http"

Limitations

在 Azure Databricks 中使用 Spark 数据源 API 时,以下限制适用。

  • 捆绑数据源的 Spark 选项仅限于 querydbtable 以及一小部分特定于连接器的选项。
  • 自定义第三方 Spark JAR 只能安装在专用群集上。 对于无服务器或标准群集,请使用捆绑连接器或 JDBC Unity 目录连接。
  • PySpark 自定义数据源需要 Databricks Runtime 15.4 LTS 或更高版本。