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
相关产品推荐
相关产品推荐

