从 Snowflake 读取和写入数据

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
userpassword sfUsersfPassword
databaseschemawarehouse sfDatabasesfSchemasfWarehouse
preactionspostactions 无服务器上不支持。 将逻辑迁移到 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 中,INTEGERDECIMAL 在语义上是等效的(请参阅 Snowflake 数值数据类型)。

为什么 Snowflak 表架构中的字段总是大写的?

默认情况下,Snowflake 使用大写字段,这意味着 Snowflake 将表架构转换为大写。