Azure Databricks Databricks Runtime 中具有 Snowflake 连接器,用于从 Snowflake 读取和写入数据。
重要
旧查询联合文档已停用,可能不会更新。 此内容中提到的配置未经 Databricks 正式认可或测试。 如果 Lakehouse 联邦 支持源数据库,Databricks 建议改用它。
在 Azure Databricks 中查询 Snowflake 表
可以配置与 Snowflake 的连接,然后查询数据。 在开始之前,请检查群集运行的 Databricks Runtime 版本。 以下代码包括 Python、SQL 和 Scala 中的示例语法。
Python
# The following example applies to Databricks Runtime 11.3 LTS and above.
snowflake_table = (spark.read
.format("snowflake")
.option("host", "hostname")
.option("port", "port") # Optional - will use default port 443 if not specified.
.option("user", "username")
.option("password", "password")
.option("sfWarehouse", "warehouse_name")
.option("database", "database_name")
.option("schema", "schema_name") # Optional - will use default schema "public" if not specified.
.option("dbtable", "table_name")
.load()
)
# The following example applies to Databricks Runtime 10.4 and below.
snowflake_table = (spark.read
.format("snowflake")
.option("dbtable", table_name)
.option("sfUrl", database_host_url)
.option("sfUser", username)
.option("sfPassword", password)
.option("sfDatabase", database_name)
.option("sfSchema", schema_name)
.option("sfWarehouse", warehouse_name)
.load()
)
SQL
/* The following example applies to Databricks Runtime 11.3 LTS and above. */
DROP TABLE IF EXISTS snowflake_table;
CREATE TABLE snowflake_table
USING snowflake
OPTIONS (
host '<hostname>',
port '<port>', /* Optional - will use default port 443 if not specified. */
user '<username>',
password '<password>',
sfWarehouse '<warehouse_name>',
database '<database-name>',
schema '<schema-name>', /* Optional - will use default schema "public" if not specified. */
dbtable '<table-name>'
);
SELECT * FROM snowflake_table;
/* The following example applies to Databricks Runtime 10.4 LTS and below. */
DROP TABLE IF EXISTS snowflake_table;
CREATE TABLE snowflake_table
USING snowflake
OPTIONS (
dbtable '<table-name>',
sfUrl '<database-host-url>',
sfUser '<username>',
sfPassword '<password>',
sfDatabase '<database-name>',
sfSchema '<schema-name>',
sfWarehouse '<warehouse-name>'
);
SELECT * FROM snowflake_table;
Scala(编程语言)
# The following example applies to Databricks Runtime 11.3 LTS and above.
val snowflake_table = spark.read
.format("snowflake")
.option("host", "hostname")
.option("port", "port") /* Optional - will use default port 443 if not specified. */
.option("user", "username")
.option("password", "password")
.option("sfWarehouse", "warehouse_name")
.option("database", "database_name")
.option("schema", "schema_name") /* Optional - will use default schema "public" if not specified. */
.option("dbtable", "table_name")
.load()
# The following example applies to Databricks Runtime 10.4 and below.
val snowflake_table = spark.read
.format("snowflake")
.option("dbtable", table_name)
.option("sfUrl", database_host_url)
.option("sfUser", username)
.option("sfPassword", password)
.option("sfDatabase", database_name)
.option("sfSchema", schema_name)
.option("sfWarehouse", warehouse_name)
.load()
将数据写入 Snowflake
可以使用与 snowflake 相同的数据源格式 df.write将 Spark 数据帧写入 Snowflake 表。 您还可以针对 Snowflake 支持的表发出 INSERT INTO 和 CTAS 语句。
Python
sf_options = {
"host": "<hostname>",
"sfDatabase": "<database-name>",
"sfSchema": "<schema-name>",
"sfWarehouse": "<warehouse-name>",
"sfRole": "<role-name>",
"sfUser": "<username>",
"sfPassword": "<password>",
"dbtable": "<table-name>",
}
(df.write
.format("snowflake")
.options(**sf_options)
.mode("append")
.save())
SQL
CREATE TABLE snowflake_target
USING snowflake
OPTIONS (
host '<hostname>',
sfUser '<username>',
sfPassword '<password>',
sfDatabase '<database-name>',
sfSchema '<schema-name>',
sfWarehouse '<warehouse-name>',
sfRole '<role-name>',
dbtable '<table-name>'
) AS SELECT * FROM source_view;
在无服务器计算上写入
在无服务器 Spark 和 Databricks SQL 上,Snowflake 连接器仅支持下表中列出的选项。
| 经典 Databricks Runtime 选项 | 无服务器对应项 |
|---|---|
sfUrl |
host (加可选 port) |
user、password |
sfUser、sfPassword |
database、schema、warehouse |
sfDatabase、sfSchema、sfWarehouse |
preactions、postactions |
无服务器上不支持。 将逻辑迁移到 Snowflake 存储过程、任务或定时作业中。 |
tempdir |
不支持无服务器环境。 |
笔记本示例:适用于 Spark 的 Snowflake 连接器
以下笔记本提供有关如何将数据写入到 Snowflake 或从 Snowflake 读取数据的简单示例。 有关详细信息,请参阅适用于 Spark 的 Snowflake 连接器。
小窍门
使用笔记本中演示的机密,避免在笔记本中公开 Snowflake 用户名和密码。
Snowflake Python 笔记本
笔记本示例:将模型训练结果保存到 Snowflake
以下笔记本介绍如何使用适用于 Spark 的 Snowflake 连接器的最佳做法。 它将数据写入 Snowflake,对一些基本数据操作使用 Snowflake,在 Azure Databricks 中训练机器学习 (ML) 模型,并将结果写回到 Snowflake。
在 Snowflake 笔记本中存储 ML 训练结果
常见问题 (FAQ)
为什么我的 Spark 数据帧列在 Snowflake 中的顺序不相同?
Spark 的 Snowflake 连接器不尊重正在写入的表的列顺序。 必须明确指定 DataFrame 与 Snowflake 列之间的映射关系。 若要指定此映射,请使用 columnmap 参数。
为什么向 Snowflake 写入的 INTEGER 数据作为 DECIMAL 读回?
Snowflak 将所有 INTEGER 类型都表示为 NUMBER,这可能会导致在将数据写入到 Snowflak 并从中读取数据时数据类型发生更改。 例如,Snowflake 可以在写入时将 INTEGER 数据转换为 DECIMAL,因为在 Snowflake 中,INTEGER 和 DECIMAL 在语义上是等效的(请参阅 Snowflake 数值数据类型)。
为什么 Snowflak 表架构中的字段总是大写的?
默认情况下,Snowflake 使用大写字段,这意味着 Snowflake 将表架构转换为大写。