Snowflake生成的Parquet文件无法被外部工具读取的问题
Snowflake导出Parquet外部读取失败的解决方法
问题原因
Snowpark将DataFrame写入Parquet格式到Stage时,会自动将数据拆分为多个Parquet分片文件(命名格式类似part-00000-xxxx.snappy.parquet),并存储在你指定的report1子目录下。你的原代码错误地读取了Session Stage根目录下的文件(大概率不是实际的Parquet数据文件),导致外部工具无法识别合法的Parquet魔数,抛出格式错误。
此外,直接将sfcfs.open()返回的文件对象传入pandas.read_parquet,会存在二进制流处理的兼容性问题,正确的做法是结合文件系统对象直接读取文件路径。
修复代码
from snowflake.snowpark import Session import pandas from snowflake.ml.fileset import sfcfs import pyarrow.parquet as pq # 替换为你的Snowflake连接配置 connection_parameters = {} # 建立Snowpark会话 snowpark_session = Session.builder.configs(connection_parameters).create() # 创建测试数据并写入Parquet到Session Stage的report1目录 df = snowpark_session.createDataFrame(pandas.DataFrame({'a': [1,2,3]})) stage_root = snowpark_session.get_session_stage() target_dir = f'{stage_root}/report1' df.write.parquet(target_dir, header=True, overwrite=True) # 初始化Snowflake文件系统 fs = sfcfs.SFFileSystem(snowpark_session=snowpark_session) # 获取report1目录下的所有Parquet文件 parquet_files = [file_path for file_path in fs.ls(target_dir) if file_path.endswith('.parquet')] # 方法1:使用PyArrow读取并转换为Pandas DataFrame parquet_table = pq.read_table(parquet_files[0], filesystem=fs) result_df = parquet_table.to_pandas() # 方法2:直接使用Pandas读取(需指定filesystem参数) result_df = pandas.read_parquet(parquet_files[0], filesystem=fs) # 输出结果 print(result_df)
核心修复要点
- 定位正确的文件路径:必须遍历你指定的
report1子目录,筛选出.parquet后缀的实际数据文件,而非Stage根目录下的无关文件。 - 利用文件系统对象读取:通过
filesystem参数将sfcfs对象传递给PyArrow或Pandas的读取方法,让库自行处理Snowflake Stage文件的流读取,避免手动打开文件导致的格式解析错误。
内容的提问来源于stack exchange,提问作者Ophir Yoktan
相关产品推荐
相关产品推荐

