能否通过Snowpark(含Snowflake Notebook)从SFTP/Oracle加载数据至Snowflake表?
解决方案:Snowpark Notebook 加载SFTP/Oracle数据到Snowflake表
一、通过Snowpark Notebook从SFTP加载文件到Snowflake表
Snowpark本身没有内置SFTP客户端,但可以在Notebook中借助Python的SFTP库(如paramiko)先将文件从SFTP服务器下载到Snowflake Internal Stage,再通过COPY INTO命令加载到目标表。具体步骤如下:
- 安装SFTP依赖库
在Snowpark Notebook的代码单元格中执行安装命令(需确保计算池允许安装第三方库):
!pip install paramiko
- 连接SFTP并将文件上传至Snowflake Stage
编写Snowpark Python代码,完成SFTP连接后将文件上传到Internal Stage:
import paramiko from snowflake.snowpark import Session # 初始化Snowpark会话 session = Session.builder.getOrCreate() # SFTP配置信息 sftp_host = "你的SFTP服务器地址" sftp_port = 22 sftp_username = "SFTP用户名" sftp_password = "SFTP密码" sftp_file_path = "/远程服务器文件路径/目标文件.csv" stage_name = "@你的内部Stage名称" stage_file_path = "上传后的文件名.csv" # 连接SFTP并将文件上传到Snowflake Stage with paramiko.Transport((sftp_host, sftp_port)) as transport: transport.connect(username=sftp_username, password=sftp_password) with paramiko.SFTPClient.from_transport(transport) as sftp: with sftp.open(sftp_file_path, 'rb') as remote_file: session.file.put_stream(remote_file, f"{stage_name}/{stage_file_path}", overwrite=True) # 验证Stage中的文件 print(session.sql(f"LIST {stage_name}").collect())
如果SFTP使用密钥认证,可将password参数替换为pkey=paramiko.RSAKey.from_private_key_file("/本地密钥文件路径")。
- 从Stage加载数据到Snowflake表
提前创建好与文件结构匹配的目标表,执行COPY INTO命令完成加载:
session.sql(""" COPY INTO 你的目标表名 FROM @你的内部Stage名称/上传后的文件名.csv FILE_FORMAT = (TYPE = CSV FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1) ON_ERROR = 'CONTINUE' """).collect()
二、通过Snowpark将Oracle中的数据/文件加载到Snowflake表
分两种场景处理:
场景1:Oracle数据库中的BLOB/CLOB类型文件
如果文件存储在Oracle的BLOB/CLOB字段中,可通过Snowpark的JDBC连接读取数据后写入Snowflake表:
配置Oracle JDBC驱动
将Oracle JDBC驱动(如ojdbc8.jar)上传到Snowflake Internal Stage,确保计算池能访问Oracle数据库网络。读取Oracle数据并写入Snowflake
from snowflake.snowpark import Session session = Session.builder.getOrCreate() # Oracle JDBC连接参数 oracle_jdbc_url = "jdbc:oracle:thin:@//Oracle服务器地址:1521/数据库实例名" oracle_username = "Oracle用户名" oracle_password = "Oracle密码" oracle_query = "SELECT id, 存储文件的BLOB字段 FROM Oracle表名" # 读取Oracle数据 oracle_df = session.read.jdbc( url=oracle_jdbc_url, table=f"({oracle_query})", properties={"user": oracle_username, "password": oracle_password, "driver": "oracle.jdbc.OracleDriver"} ) # 写入Snowflake目标表 oracle_df.write.mode("append").save_as_table("Snowflake目标表名")
场景2:Oracle服务器本地文件
如果文件存储在Oracle服务器的文件系统中,需先将文件迁移到Snowflake可访问的位置:
- 方式1:若Snowpark计算池与Oracle服务器网络互通,可使用Python文件操作库读取文件后上传到Snowflake Stage,再加载到表(步骤同SFTP场景)。
- 方式2:通过Oracle的
UTL_FILE工具将文件导出到云存储(如S3、Azure Blob),再使用Snowpark的session.file.put或COPY INTO命令加载到Snowflake表。
内容的提问来源于stack exchange,提问作者hari_azure
相关产品推荐
相关产品推荐

