snowflake
snowflake 插件用于查看和读取 Snowflake 数据,支持建立和关闭 Snowflake 连接,查看指定 Snowflake 数据库或架构(Schema)下当前 Snowflake 角色可见的表和视图,以及获取表结构或查询结果的列元数据。在数据读取方面,支持将 Snowflake 表或查询结果加载为 DolphinDB 内存表,也支持将表数据或查询结果分批写入已有的 DolphinDB 表。
-
暂不支持将 DolphinDB 数据写入 Snowflake。
-
如需下载大量数据,请参考 Snowflake 官方文档。
前提条件
安装插件
版本要求
DolphinDB Server 2.00.16 及更高版本、3.00.4 及更高版本,安装包类型为 Linux x86_64 或 Linux x86_64 ABI。
安装步骤
-
在 DolphinDB 客户端中使用 listRemotePlugins 函数查看可供安装的插件。
login("admin", "123456") listRemotePlugins() -
使用 installPlugin 函数安装插件。
installPlugin("snowflake") -
使用 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,不得包含协议(如
|
|
"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 请求错误的重试次数;总尝试次数最多为
|
|
"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 列名。列名规范化规则如下:
|
|
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 列名。列名规范化规则如下:
|
|
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" |
一张表,用于覆盖结果列名和转换目标类型。表必须满足以下条件:
|
|
|
"offset" |
"offset" 用于决定从第几条数据开始返回,"limit" 用于限制返回的数据条数。
|
|
|
"limit" |
||
|
"allowEmptyResult" |
BOOL 类型标量,指定是否允许返回空表,默认为 false。
|
|
|
"batchTransform" |
自定义 DolphinDB 函数,用于处理批次数据。
|
|
支持的数据类型
Snowflake 与 DolphinDB 数据类型的映射
插件根据 Snowflake SQL API 返回的 metadata 中的 type、precision 和 scale 进行类型映射,而不是仅根据建表 DDL 中使用的类型名称进行判断。
| Snowflake metadata type | 条件 | 自动映射的 DolphinDB 类型 |
|---|---|---|
|
fixed、number、decimal、numeric |
|
INT |
|
|
LONG |
|
|
|
STRING |
|
|
|
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)
