snowflake

snowflake 插件用于查看和读取 Snowflake 数据,支持建立和关闭 Snowflake 连接,查看指定 Snowflake 数据库或架构(Schema)下当前 Snowflake 角色可见的表和视图,以及获取表结构或查询结果的列元数据。在数据读取方面,支持将 Snowflake 表或查询结果加载为 DolphinDB 内存表,也支持将表数据或查询结果分批写入已有的 DolphinDB 表。

注:
  • 暂不支持将 DolphinDB 数据写入 Snowflake。

  • 如需下载大量数据,请参考 Snowflake 官方文档

前提条件

  • DolphinDB Server 所在服务器必须配置可用的 CA 证书以及 RSA 私钥文件

  • DolphinDB Server 所在服务器的系统时间应保持准确,否则 JWT 可能因签发时间或过期时间校验失败而导致认证失败。

  • 已创建好 Snowflake 组织账号用户仓库数据库架构

  • 生成与 DolphinDB Server 所在服务器中 RSA 私钥匹配的公钥,将公钥配置给 Snowflake 用户。

    ALTER USER [USERNAME] SET RSA_PUBLIC_KEY = 'RSA public key'

安装插件

版本要求

DolphinDB Server 2.00.16 及更高版本、3.00.4 及更高版本,安装包类型为 Linux x86_64 或 Linux x86_64 ABI

安装步骤

  1. 在 DolphinDB 客户端中使用 listRemotePlugins 函数查看可供安装的插件。

    login("admin", "123456")
    listRemotePlugins()
  2. 使用 installPlugin 函数安装插件。

    installPlugin("snowflake")
  3. 使用 loadPlugin 函数加载插件。

    loadPlugin("snowflake")

接口说明

connect

语法

snowflake::connect(config)

详情

建立与 Snowflake 数据库的连接,并返回一个连接句柄。

注:

该接口通过 JWT 完成 Snowflake 密钥对身份认证。连接时,插件从 "privateKeyPath" 指定的本地私钥文件路径中读取 RSA 私钥,并基于 "account" 和 "user" 生成 JWT。插件通过 Authorization 请求头将 JWT 发送至 Snowflake,Snowflake 使用预先配置在对应用户上的 RSA 公钥验证 JWT,验证通过后完成身份认证。

参数

config 字典(Dictionary<STRING, ANY>),指定连接配置,支持以下键值对:

是否必填

"account"

STRING 类型标量,指定 Snowflake 组织中的账号标识符(不指定区域)。标识符格式为 "organization-account",例如 "WJCYPAV-TD40511"。

"user"

STRING 类型标量,指定 Snowflake 账号中的用户名。

"privateKeyPath"

STRING 类型标量,指定 DolohinDB Server 中的 RSA 私钥文件路径。

"warehouse"

STRING 类型标量,指定提供计算资源的仓库。默认为 "warehouse"。

"database"

STRING 类型标量,指定要访问的数据库。默认为 "database"。

"schema"

STRING 类型标量,指定数据库中的架构。默认为 "schema"。

"role"

STRING 类型标量,指定连接 Snowflake 时使用的角色。默认为 Snowflake 用户自身配置的角色。

"options"

字典(Dictionary<STRING, ANY>),指定主机名、JWT、轮询间隔和超时时间等配置。详情见下表。

options 字典支持以下键值对:

"host"

STRING 类型标量,指定 Snowflake 主机名,仅允许填写 hostname,不得包含协议(如 https://)、端口、路径或 URL 参数(包括 query 和 fragment)。默认将账号标识符转为小写并拼接 .snowflakecomputing.com 作为主机名,例如 "myorg-myaccount.snowflakecomputing.com"。

"privateKeyPassphrase"

STRING 类型标量,指定用于解密 RSA 私钥的口令。仅在本地 RSA 私钥加密时配置,默认为空字符串。

"jwtLifetimeSeconds"

取值范围为 [1, 3600] 的整数标量,指定 JWT 的有效期,单位为秒。默认为 3600。

"jwtRefreshMarginSeconds"

取值范围为 [0, 3600] 的整数标量,指定 JWT 提前刷新时间,单位为秒。当 JWT 距离过期时间小于该值时,插件会刷新 JWT。该值必须小于 jwtLifetimeSeconds。默认为 300。

"connectTimeoutSeconds"

正整数标量,指定 libcurl 建立 TCP/TLS 连接的超时时间,单位为秒。默认为 10。

"requestTimeoutSeconds"

正整数标量,指定 Snowflake 语句请求的超时时间,单位为秒。该参数仅用于设置请求体中的超时时间,不表示 libcurl 整个 HTTP 请求的总超时时间。默认为 60。

"pollIntervalMillis"

正整数标量,指定 SQL API 返回 202 后的轮询间隔,单位为毫秒。默认为 500。

"maxPollSeconds"

正整数标量,指定异步语句的最大轮询时间,单位为秒。默认为 3600。

"maxRetry"

非负整数标量,指定 429、500、502、503、504 和 libcurl 请求错误的重试次数;总尝试次数最多为 maxRetry + 1。默认为 3。

"proxy"

STRING 类型标量,指定传给 libcurl 的 HTTP/HTTPS 代理。默认为空字符串,表示不设置代理。

返回值

一个连接句柄。

示例

// 基本配置
config = dict(STRING, ANY)
config["account"] = "myorg-myaccount"
config["user"] = "DDB_READER"
config["privateKeyPath"] = "/opt/dolphindb/keys/snowflake_rsa_key.p8"
config["warehouse"] = "DDB_WH"
config["database"] = "SALES_DB"
config["schema"] = "PUBLIC"
config["role"] = "DDB_READER_ROLE"

// 通过 options 字典设置私钥口令和超时时间
connOptions = dict(STRING, ANY)
connOptions["privateKeyPassphrase"] = "<PRIVATE_KEY_PASSPHRASE>"
connOptions["connectTimeoutSeconds"] = 20
connOptions["requestTimeoutSeconds"] = 120
connOptions["pollIntervalMillis"] = 1000
connOptions["maxPollSeconds"] = 1800
connOptions["maxRetry"] = 5
config["options"] = connOptions

conn = snowflake::connect(config)

close

语法

snowflake::close(conn)

详情

关闭指定的连接。连接关闭后,若将对应的句柄传给任意插件接口,运行接口时将报错:Invalid connection object.

参数

connconnect 创建的连接句柄。

返回值

字符串 "Connection is closed."

showTables

语法

snowflake::showTables(conn, [database], [schema])

详情

查看当前 Snowflake 用户可见的表或视图。

参数

connconnect 创建的连接句柄。

database STRING 类型标量,指定要查询的数据库。默认为创建对应连接句柄时指定的数据库。

schema STRING 类型标量,指定要查询的架构。默认为创建对应连接句柄时指定的架构。

返回值

一张表,包含以下列,没有可见对象时返回空表:

列名 数据类型 说明

database

STRING

Snowflake 数据库。

schema

STRING

Snowflake 数据库的架构。

name

STRING

Snowflake 表或视图名称。

kind

STRING

Snowflake 数据库对象的类型,例如 BASE TABLE 或 VIEW。

示例

// 查询连接默认的数据库和架构
objects = snowflake::showTables(conn)

// 指定数据库,架构仍使用连接默认值
objects = snowflake::showTables(conn, "ARCHIVE_DB")

// 同时指定数据库和架构
objects = snowflake::showTables(conn, "ARCHIVE_DB", "REPORTING")

// 只指定架构,可把 database 参数设置为 NULL
objects = snowflake::showTables(conn, NULL, "REPORTING")

extractTableSchema

语法

snowflake::extractTableSchema(conn, sourceTableName)

详情

获取 Snowflake 表的结构。

参数

connconnect 创建的连接句柄。

sourceTableName STRING 类型标量,指定 Snowflake 表的名称,支持 1-3 段表名形式:TABLE_NAME、SCHEMA_NAME.TABLE_NAME、DATABASE_NAME.SCHEMA_NAME.TABLE_NAME。

返回值

一张表,包含以下列:

列名 数据类型 说明

name

STRING

规范化后的 DolphinDB 列名。列名规范化规则如下:

  • 空列名转换为 col_1、col_2 等。

  • 非字母、数字或下划线的字符替换为下划线。

  • 如果列名的首字符不是字母或下划线,则添加 col_ 前缀

  • 对规范化后重名的列,依次添加 _2、_3 等后缀以确保列名唯一。

type

STRING

根据 Snowflake 数据类型自动映射得到的 DolphinDB 数据类型。

snowflakeType

STRING

Snowflake SQL API 返回的原始数据类型。

precision

INT

Snowflake metadata 中的精度。如果 metadata 中未提供,则返回 -1。

scale

INT

Snowflake metadata 中的小数位数。如果 metadata 中未提供,则返回 0。

nullable

BOOL

表示该列是否允许为 NULL,取值来源于 Snowflake metadata。

extractQuerySchema

语法

snowflake::extractQuerySchema(conn, sql)

详情

对 Snowflake 表执行指定的 SELECT 语句,获取查询结果的结构。

参数

connconnect 创建的连接句柄。

sql STRING 类型标量,指定查询 SQL。插件会删除 SQL 尾部的分号,并将查询包装为类似以下的语句:

SELECT *
FROM (<sql>) AS SNOWFLAKE_PLUGIN_SCHEMA
WHERE 1 = 0

因此,sql 必须是可作为子查询使用的 SELECT 语句。DDL 或 DML 语句不适合传给该接口。

返回值

一张表,包含以下列:

列名 数据类型 说明

name

STRING

规范化后的 DolphinDB 列名。列名规范化规则如下:

  • 空列名转换为 col_1、col_2 等。

  • 非字母、数字或下划线的字符替换为下划线。

  • 如果列名的首字符不是字母或下划线,则添加 col_ 前缀

  • 对规范化后重名的列,依次添加 _2、_3 等后缀以确保列名唯一。

type

STRING

根据 Snowflake 数据类型自动映射得到的 DolphinDB 数据类型。

snowflakeType

STRING

Snowflake SQL API 返回的原始数据类型。

precision

INT

Snowflake metadata 中的精度。如果 metadata 中未提供,则返回 -1。

scale

INT

Snowflake metadata 中的小数位数。如果 metadata 中未提供,则返回 0。

nullable

BOOL

表示该列是否允许为 NULL,取值来源于 Snowflake metadata。

loadTable

语法

snowflake::loadTable(conn, sourceTableName, [options])

详情

加载 Snowflake 表数据并返回为 DolphinDB 内存表。

参数

connconnect 创建的连接句柄。

sourceTableName STRING 类型标量,指定 Snowflake 表的名称,支持 1-3 段表名形式:TABLE_NAME、SCHEMA_NAME.TABLE_NAME、DATABASE_NAME.SCHEMA_NAME.TABLE_NAME。

options 可选参数,字典(Dictionary<STRING, ANY>),支持 "columnSchema"、"offset"、"limit"、"allowEmptyResult" 键,详情参见数据查询接口的公共参数

返回值

DolphinDB 内存表。

示例

// 不指定额外配置 options
orders = snowflake::loadTable(conn, "SALES_DB.PUBLIC.ORDERS")

// 分页返回数据
options = dict(STRING, ANY)
options["offset"] = 1000
options["limit"] = 500
options["allowEmptyResult"] = true

ordersPage = snowflake::loadTable(
    conn,
    "SALES_DB.PUBLIC.ORDERS",
    options
)

// 自定义列类型
columnSchema = table(
    ["order_id", "customer", "amount"] as name,
    ["LONG", "STRING", "DOUBLE"] as type
)

options = dict(STRING, ANY)
options["columnSchema"] = columnSchema

orders = snowflake::loadTable(
    conn,
    "SALES_DB.PUBLIC.ORDER_SUMMARY",
    options
)

query

语法

snowflake::query(conn, sql, [options])

详情

对 Snowflake 表执行指定的 SELECT 语句,并将查询结果返回为 DolphinDB 内存表。

参数

connconnect 创建的连接句柄。

sql STRING 类型标量,指定查询 SQL。插件会删除 SQL 尾部的分号,并将查询包装为类似以下的语句:

SELECT *
FROM (<sql>) AS SNOWFLAKE_PLUGIN_SCHEMA
WHERE 1 = 0

因此,sql 必须是可作为子查询使用的 SELECT 语句。DDL 或 DML 语句不适合传给该接口。

options 可选参数,字典(Dictionary<STRING, ANY>),支持 "columnSchema"、"allowEmptyResult" 键,详情参见数据查询接口的公共参数

注:

query 不支持 options 中的 "offset" 和 "limit";应直接在 SQL 中写 LIMIT、OFFSET 和 ORDER BY。

返回值

DolphinDB 内存表。

示例

// 不指定额外配置 options
sql = "SELECT ORDER_ID, AMOUNT FROM SALES_DB.PUBLIC.ORDERS WHERE AMOUNT > 100 ORDER BY ORDER_ID"
result = snowflake::query(conn, sql)

// 分页返回数据
sql = "SELECT ORDER_ID, AMOUNT FROM SALES_DB.PUBLIC.ORDERS ORDER BY ORDER_ID LIMIT 500 OFFSET 1000"

options = dict(STRING, ANY)
options["allowEmptyResult"] = true

result = snowflake::query(conn, sql, options)

loadTableInto

语法

snowflake::loadTableInto(conn, sourceTableName, targetTable, [options])

详情

加载 Snowflake 表数据并将其写入 DolphinDB 表。

注:

写入不保证事务原子性。数据按 partition 依次写入,部分 partition 写入成功后,后续 partition 即使失败,已写入的数据也不会自动回滚。如需保证原子性,应先将数据写入临时表,完成数据校验后,再由业务脚本将临时表切换为正式表。

参数

connconnect 创建的连接句柄。

sourceTableName STRING 类型标量,指定 Snowflake 表的名称,支持 1-3 段表名形式:TABLE_NAME、SCHEMA_NAME.TABLE_NAME、DATABASE_NAME.SCHEMA_NAME.TABLE_NAME。

targetTable 一个表对象,可以是 DolphinDB 内存表或分布式表。

options 可选参数,字典(Dictionary<STRING, ANY>),支持 "columnSchema"、"offset"、"limit"、"batchTransform" 键,详情参见数据查询接口的公共参数

返回值

示例

// 写入 DolphinDB 内存表
target = table(
    100000:0,
    `order_id`customer`amount,
    [LONG, STRING, DOUBLE]
)

writtenRows = snowflake::loadTableInto(
    conn,
    "SALES_DB.PUBLIC.ORDER_SUMMARY",
    target
)

// 限制读取行数
options = dict(STRING, ANY)
options["limit"] = 10000

writtenRows = snowflake::loadTableInto(
    conn,
    "SALES_DB.PUBLIC.ORDER_SUMMARY",
    target,
    options
)

queryInto

语法

snowflake::queryInto(conn, sql, targetTable, [options])

详情

对 Snowflake 表执行指定的 SELECT 语句,将查询结果写入 DolphinDB 表。

注:

写入不保证事务原子性。数据按 partition 依次写入,部分 partition 写入成功后,后续 partition 即使失败,已写入的数据也不会自动回滚。如需保证原子性,应先将数据写入临时表,完成数据校验后,再由业务脚本将临时表切换为正式表。

参数

connconnect 创建的连接句柄。

sql STRING 类型标量,指定查询 SQL。插件会删除 SQL 尾部的分号,并将查询包装为类似以下的语句:

SELECT *
FROM (<sql>) AS SNOWFLAKE_PLUGIN_SCHEMA
WHERE 1 = 0

因此,sql 必须是可作为子查询使用的 SELECT 语句。DDL 或 DML 语句不适合传给该接口。

targetTable 一个表对象,可以是 DolphinDB 内存表或分布式表。

options 可选参数,字典(Dictionary<STRING, ANY>),支持 "columnSchema"、"batchTransform" 键,详情参见数据查询接口的公共参数

注:

queryInto 不支持 options 中的 "offset" 和 "limit";应直接在 SQL 中写 LIMIT、OFFSET 和 ORDER BY。

返回值

示例

// 写入 DolphinDB 内存表
target = table(
    100000:0,
    `order_id`customer`amount,
    [LONG, STRING, DOUBLE]
)

sql = "SELECT ORDER_ID, CUSTOMER, AMOUNT FROM SALES_DB.PUBLIC.ORDERS ORDER BY ORDER_ID"
writtenRows = snowflake::queryInto(conn, sql, target)

// 写入 DolphinDB 分布式表
target = loadTable("dfs://sales", `orders)

sql = "SELECT ORDER_ID, CUSTOMER, AMOUNT FROM SALES_DB.PUBLIC.ORDERS"
writtenRows = snowflake::queryInto(conn, sql, target)

数据查询接口的公共参数

loadTablequeryloadTableIntoqueryInto 都支持字典形式的 options 参数,其支持的键值对如下表。每个接口支持的键不完全相同,传入不支持的键会导致接口调用报错。

支持该键的接口

"columnSchema"

一张表,用于覆盖结果列名和转换目标类型。表必须满足以下条件:

  • 至少有两列且插件只读取前两列。

  • 第一列是列名,类型必须为 STRING 或 SYMBOL;列名不能为空。

  • 第二列是类型名,类型必须为 STRING 或 SYMBOL;类型名大小写不敏感,可以带 DT_ 前缀。

  • 行数必须与 Snowflake 结果列数完全一致。

  • 当前只支持将 Snowflake 数据转换为以下类型:BOOL、INT、LONG、DOUBLE、STRING 和 SYMBOL。如果无法强制转换将报错。

loadTablequeryloadTableIntoqueryInto

"offset"

"offset" 用于决定从第几条数据开始返回,"limit" 用于限制返回的数据条数。

  • "offset" 和 "limit" 必须是非负整数标量。

  • "offset" > 0 时必须同时指定 "limit"。

  • 设置 "limit" 后生成 LIMIT <limit> 子句;"offset" > 0 时生成 OFFSET <offset> 子句。

  • 指定 "offset" 和 "limit" 不保证表的读取顺序;按指定顺序稳定分页应使用 query/queryInto 并在 SQL 中使用 ORDER BY。

loadTableloadTableInto

"limit"

"allowEmptyResult"

BOOL 类型标量,指定是否允许返回空表,默认为 false。

  • 默认情况下,0 行结果会导致报错: Snowflake query returned an empty table.

  • 设置为 true 时支持返回保留列结构的 0 行 DolphinDB 表。

loadTablequery

"batchTransform"

自定义 DolphinDB 函数,用于处理批次数据。

  • 函数的输入和输出都必须是表。

  • 返回 0 行表时,本批次不调用 tableInsert

  • 返回的非空表必须与 targetTable 的列数及对应的 DolphinDB 数据类型完全一致。

  • 函数抛出的异常会直接返回给调用者。

loadTableIntoqueryInto

支持的数据类型

Snowflake 与 DolphinDB 数据类型的映射

插件根据 Snowflake SQL API 返回的 metadata 中的 type、precision 和 scale 进行类型映射,而不是仅根据建表 DDL 中使用的类型名称进行判断。

Snowflake metadata type 条件 自动映射的 DolphinDB 类型

fixed、number、decimal、numeric

scale = 00 <= precision <= 9

INT

scale = 0precision < 09 < precision <= 18

LONG

scale = 0precision > 18

STRING

scale != 0

DOUBLE

real、float、double

任意

DOUBLE

boolean、bool

任意

BOOL

其他所有 metadata type

任意

STRING

以下 Snowflake 数据类型当前通常映射为 STRING:

  • CHAR、VARCHAR、TEXT

  • DATE、TIME

  • TIMESTAMP_NTZ、TIMESTAMP_LTZ、TIMESTAMP_TZ

  • BINARY、VARBINARY

  • ARRAY、OBJECT、VARIANT

  • GEOGRAPHY、GEOMETRY、VECTOR

  • 其他未在上述映射规则中列出的类型

对于映射为 STRING 的类型,返回值为 Snowflake SQL API 返回值的文本表示。特别是 DATE、TIME 和 TIMESTAMP 类型,当前不应假定其返回值为 DolphinDB 原生时间类型。

columnSchema 支持的目标类型

columnSchema 支持以下目标类型及其别名。类型名大小写不敏感。

类型名 可用别名 接受的非 NULL 文本

BOOL

BOOLEAN、DT_BOOL、DT_BOOLEAN

true、false、1、0,大小写不敏感

INT

DT_INT

完整的十进制整数,且值在 DolphinDB INT 的可表示范围内

LONG

DT_LONG

完整的十进制整数,且值在 DolphinDB LONG 的可表示范围内

DOUBLE

DT_DOUBLE

可完整解析为 DolphinDB DOUBLE 的文本

STRING

DT_STRING

任意非空文本

SYMBOL

DT_SYMBOL

任意非空文本

非 NULL 文本是指 Snowflake SQL API 返回的非 NULL 字段值的文本表示。

如果 columnSchema 中指定了上述列表之外的类型,接口将返回错误:schema type ... is not supported yet

NULL 值与类型转换

插件请求 Snowflake SQL API 时设置 nullable=true,并按以下规则处理返回值:

  • Snowflake SQL NULL:以 JSON null 返回,并转换为对应 DolphinDB 类型的 NULL。

  • 字符串 "null":作为普通字符串处理,不视为 NULL。

  • 非 NULL 空字符串:DolphinDB 使用空字符串表示 NULL。因此,Snowflake 中的非 NULL 空字符串如果直接转换会被误认为 NULL,插件将阻止转换并报错。

  • 非 NULL 数值:如果 Snowflake 返回的非 NULL 数值等于 DolphinDB 对应类型的 NULL 表示值,直接转换会被误认为 NULL。为避免数据错误,插件将阻止转换并报错。例如,Snowflake 返回的 INT 最小值、LONG 最小值,或等于 DolphinDB DOUBLE NULL 表示值的数值,均无法作为非 NULL 值转换。

  • 数值转换:数值文本必须能够被完整解析为目标类型。包含多余字符或超出目标类型可表示范围时,将报转换错误。

  • NUMBER/DECIMAL:当 scale != 0 时,默认映射为 DOUBLE,可能产生浮点精度损失。如需完整保留原始数值文本,可通过 columnSchema 将该列指定为 STRING 类型。

完整示例

例 1:以下示例演示加载插件、建立连接、查看对象、提取结构、查询数据、写入目标表和关闭连接的完整流程。

// 1. 加载插件
loadPlugin("snowflake")

// 2. 创建连接
config = dict(STRING, ANY)
config["account"] = "myorg-myaccount"
config["user"] = "DDB_READER"
config["privateKeyPath"] = "/opt/dolphindb/keys/snowflake_rsa_key.p8"
config["warehouse"] = "DDB_WH"

config["database"] = "SALES_DB"
config["schema"] = "PUBLIC"
config["role"] = "DDB_READER_ROLE"

conn = snowflake::connect(config)

// 3. 查看可见对象
objects = snowflake::showTables(conn)

// 4. 检查查询结果结构
sql = "SELECT ORDER_ID, CUSTOMER, AMOUNT FROM SALES_DB.PUBLIC.ORDERS ORDER BY ORDER_ID"
sourceSchema = snowflake::extractQuerySchema(conn, sql)

// 5. 小结果直接返回为内存表
options = dict(STRING, ANY)
options["allowEmptyResult"] = true
preview = snowflake::query(conn, sql + " LIMIT 100", options)

// 6. 大结果分批写入已有目标表
target = table(
    100000:0,
    `order_id`customer`amount,
    [LONG, STRING, DOUBLE]
)

columnSchema = table(
    ["order_id", "customer", "amount"] as name,
    ["LONG", "STRING", "DOUBLE"] as type
)

writeOptions = dict(STRING, ANY)
writeOptions["columnSchema"] = columnSchema

writtenRows = snowflake::queryInto(conn, sql, target, writeOptions)

// 7. 使用完毕后关闭连接
snowflake::close(conn)

例 2:使用 columnSchema 显示指定目标类型。

columnSchema = table(
    ["order_id", "customer", "amount"] as name,
    ["LONG", "STRING", "DOUBLE"] as type
)

options = dict(STRING, ANY)
options["columnSchema"] = columnSchema

target = table(
    100000:0,
    `order_id`customer`amount,
    [LONG, STRING, DOUBLE]
)

sql = "SELECT ORDER_ID, CUSTOMER, AMOUNT FROM SALES_DB.PUBLIC.ORDERS"
writtenRows = snowflake::queryInto(conn, sql, target, options)