You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Snowflake存储过程下载CSV至内部阶段报错:流不可查找

问题:Snowflake存储过程put_stream写入内部阶段报错io.UnsupportedOperation: File or stream is not seekable

我希望通过程序化方式从某公共网站下载数据,该网站提供静态URL用于下载每周更新的CSV文件。方案是使用具备外部网络访问权限的Snowflake存储过程将文件下载至内部阶段,再基于这些文件构建表或视图,最后通过任务实现每周自动下载新文件。

编写的存储过程能够访问文件流(外部访问配置有效),但写入内部阶段时失败:

create or replace procedure download_data(file_url string, dest_stage string, dest_filename string)
returns variant
language python
runtime_version = 3.11
packages = ('snowflake-snowpark-python', 'urllib3')
external_access_integrations = (my_access_integration)
handler = 'main'
execute as caller
as
$$

from snowflake.snowpark import FileOperation

import urllib3
import os


def main(session, file_url, dest_stage, dest_filename):
    
    http = urllib3.PoolManager()
    resp = http.request('GET', file_url, preload_content = False,)

    print(resp.read(10)) # this works!
    
    sf = FileOperation(session)
    sf.put_stream(
        input_stream = resp,
        stage_location = os.path.join(dest_stage, dest_filename),
        auto_compress = False,
    )

    return f'File from ({file_url}) landed in stage {dest_stage} with name {dest_filename}'
    
$$
;

报错信息:

io.UnsupportedOperation: File or stream is not seekable.

请问是对文件流的使用方式有误,还是put_stream方法暂不支持该行为?


解决方案

问题根源:urllib3返回的响应流不支持随机访问(不可seek),而Snowpark的put_stream方法要求输入流必须具备seek能力(需要获取流长度、重新定位指针等操作)。

以下两种方法可以解决问题:

方法1:将响应内容转存到可seek的内存流(BytesIO)

通过io.BytesIO将urllib3的响应内容读取到内存中,生成支持seek操作的流,再传递给put_stream:

create or replace procedure download_data(file_url string, dest_stage string, dest_filename string)
returns variant
language python
runtime_version = 3.11
packages = ('snowflake-snowpark-python', 'urllib3', 'io')
external_access_integrations = (my_access_integration)
handler = 'main'
execute as caller
as
$$

from snowflake.snowpark import FileOperation
from io import BytesIO

import urllib3
import os


def main(session, file_url, dest_stage, dest_filename):
    
    http = urllib3.PoolManager()
    resp = http.request('GET', file_url, preload_content = False)

    # 将响应流读取到BytesIO(支持seek的内存流)
    stream = BytesIO(resp.read())
    
    sf = FileOperation(session)
    sf.put_stream(
        input_stream = stream,
        stage_location = os.path.join(dest_stage, dest_filename),
        auto_compress = False,
    )

    return f'File from ({file_url}) landed in stage {dest_stage} with name {dest_filename}'
    
$$
;

方法2:改用requests库(响应流默认支持seek)

requests库的响应原生流resp.raw支持seek操作,无需额外转换,直接传入put_stream即可:

create or replace procedure download_data(file_url string, dest_stage string, dest_filename string)
returns variant
language python
runtime_version = 3.11
packages = ('snowflake-snowpark-python', 'requests')
external_access_integrations = (my_access_integration)
handler = 'main'
execute as caller
as
$$

from snowflake.snowpark import FileOperation

import requests
import os


def main(session, file_url, dest_stage, dest_filename):
    
    resp = requests.get(file_url, stream=True)
    resp.raw.decode_content = True  # 处理服务器返回的gzip压缩内容
    
    sf = FileOperation(session)
    sf.put_stream(
        input_stream = resp.raw,
        stage_location = os.path.join(dest_stage, dest_filename),
        auto_compress = False,
    )

    return f'File from ({file_url}) landed in stage {dest_stage} with name {dest_filename}'
    
$$
;

内容的提问来源于stack exchange,提问作者Le Chase

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 09:23:32