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

使用transforms.api筛选含ASGK的CSV文件时代码报错求助

CSV文件筛选与读取错误排查及解决

问题背景

数据集包含多个CSV文件,需仅处理文件名含ASGK的文件(如ASGK_2022.csv、ASGK_2023.csv),计划通过transforms.api.FileSystem.files结合正则表达式筛选,但两段尝试代码均报错。

错误分析

第一段代码错误原因

代码误用pd.ExcelFile读取CSV文件:

with pd.ExcelFile(f.read(), engine='openpyxl') as xlsx_path:
    pdf = pd.read_csv(xlsx_path, dtype=str, header=0)

pd.ExcelFile是专门解析Excel(.xlsx/.xls)文件的工具,而CSV是纯文本格式,并非Zip压缩包(Excel本质是Zip格式),因此触发zipfile.BadZipFile错误。

第二段代码错误原因

代码错误使用olefile.OleFileIO处理CSV文件:

ole = olefile.OleFileIO(f.read())

olefile用于解析OLE格式文件(如旧版.xls),完全不适用于CSV文本文件。且OleFileIO接受文件路径或文件对象,传入f.read()返回的二进制内容会导致路径解析错误,触发FileNotFoundError。

解决方案

直接针对CSV文件读取处理,无需引入Excel/OLE相关库,以下是两种可行实现方式:

方式1:Pandas读取后转Spark DataFrame

from pyspark.sql import functions as F
from transforms.api import transform, Input, Output
import pandas as pd
import json


@transform(
    output_df=Output("your_output_path"),
    input_raw=Input("your_input_path"),
)
def compute(input_raw, output_df, ctx):

    def process_file(file_status):
        with input_raw.filesystem().open(file_status.path, 'r', encoding='utf-8') as f:
            # 直接读取CSV文件
            pdf = pd.read_csv(f, dtype=str, header=0)
            pdf.columns = pdf.columns.str.lower()
            
            for row in pdf.to_dict('records'):
                yield json.dumps(row, default=str)

    # 用正则筛选含ASGK的CSV文件
    files_rdd = input_raw.filesystem().files(regex=r'.*ASGK.*\.csv$').rdd.flatMap(process_file)
    spark = ctx.spark_session
    df = spark.read.json(files_rdd)
    output_df.write_dataframe(df)

方式2:直接用Spark读取(更高效)

利用Spark原生CSV读取能力,无需逐文件处理,性能更优:

from transforms.api import transform, Input, Output
from pyspark.sql import SparkSession


@transform(
    output_df=Output("your_output_path"),
    input_raw=Input("your_input_path"),
)
def compute(input_raw, output_df, ctx):
    spark = ctx.spark_session
    
    # 获取所有符合条件的文件路径
    files_df = input_raw.filesystem().files(regex=r'.*ASGK.*\.csv$')
    file_paths = [row.path for row in files_df.collect()]
    
    # 直接读取多个CSV文件
    df = spark.read.csv(
        file_paths,
        header=True,
        inferSchema=False,
        encoding='utf-8'
    )
    # 列名转小写
    df = df.toDF(*[col.lower() for col in df.columns])
    
    output_df.write_dataframe(df)

关键注意事项

  • 正则表达式r'.*ASGK.*\.csv$'确保仅匹配文件名含ASGK且以.csv结尾的文件
  • 读取CSV时指定编码(如utf-8),避免乱码问题
  • 优先使用Spark原生读取API,比逐文件用Pandas处理更高效,尤其适合大数据量场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 05:25:35