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.

参数

conn 由 connect 创建的连接句柄。

返回值

字符串 "Connection is closed."

showTables

语法

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

详情

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

参数

conn 由 connect 创建的连接句柄。

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 表的结构。

参数

conn 由 connect 创建的连接句柄。

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 语句,获取查询结果的结构。

参数

conn 由 connect 创建的连接句柄。

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 内存表。

参数

conn 由 connect 创建的连接句柄。

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 内存表。

参数

conn 由 connect 创建的连接句柄。

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 即使失败,已写入的数据也不会自动回滚。如需保证原子性,应先将数据写入临时表,完成数据校验后,再由业务脚本将临时表切换为正式表。

参数

conn 由 connect 创建的连接句柄。

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 即使失败,已写入的数据也不会自动回滚。如需保证原子性,应先将数据写入临时表,完成数据校验后,再由业务脚本将临时表切换为正式表。

参数

conn 由 connect 创建的连接句柄。

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)

数据查询接口的公共参数

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

键 值 支持该键的接口

"columnSchema"

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

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

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

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

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

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

loadTable、query、loadTableInto、queryInto

"offset"

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

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

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

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

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

loadTable、loadTableInto

"limit"

"allowEmptyResult"

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

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

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

loadTable、query

"batchTransform"

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

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

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

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

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

loadTableInto、queryInto

支持的数据类型

Snowflake 与 DolphinDB 数据类型的映射

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

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

fixed、number、decimal、numeric

scale = 0 且 0 <= precision <= 9

INT

scale = 0 且 precision < 0 或 9 < precision <= 18

LONG

scale = 0 且 precision > 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)