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

如何循环遍历表格列表,抽取首尾各1000行生成对应CSV文件

多表批量抽取首尾1000行并导出CSV实现方案

前置修复:单表抽取逻辑问题

你现有单表抽取代码中后1000行的逻辑存在错误,和前1000行逻辑完全一致,修复后的样本抽取逻辑如下:

# 修复后的首尾1000行抽取逻辑
def compute_first_last_1000(my_input):
    df_with_index = my_input.withColumn("index", monotonically_increasing_id()).orderBy("index").cache()
    # 取前1000行
    first_df = df_with_index.filter(col("index") < 1000).drop("index")
    # 取后1000行:倒序取前1000后恢复正序
    last_df = df_with_index.orderBy(col("index").desc()).limit(1000).orderBy("index").drop("index")
    return first_df.unionByName(last_df)

批量处理实现代码

直接通过循环动态生成每个表的转换任务,自动注册到Foundry上下文即可,不需要重复写120次代码:

# 导入依赖
from transforms.api import transform_df, transform, Input, Output
from pyspark.sql.functions import monotonically_increasing_id, col
import csv

# 1. 填充你的120张表名到这个列表
list_of_tables = [
    "stock",
    "sales",
    "user",
    # 补全剩余所有表名
]

# 2. 循环生成每个表的抽取+导出任务
for table_name in list_of_tables:
    # 生成样本抽取任务
    @transform_df(
        output=Output(f"foundry/sample/{table_name}_sample"),
        my_input=Input(f"foundry/input/{table_name}"),
    )
    def sample_task(my_input, table_name=table_name):
        df_with_index = my_input.withColumn("index", monotonically_increasing_id()).orderBy("index").cache()
        first_df = df_with_index.filter(col("index") < 1000).drop("index")
        last_df = df_with_index.orderBy(col("index").desc()).limit(1000).orderBy("index").drop("index")
        return first_df.unionByName(last_df)
    # 注册任务到全局上下文,Foundry可识别
    globals()[f"sample_{table_name}"] = sample_task

    # 生成CSV导出任务
    @transform(
        output=Output(f"foundry/sample/{table_name}_sample_csv"),
        sample_input=Input(f"foundry/sample/{table_name}_sample"),
    )
    def export_task(sample_input, output, table_name=table_name):
        df = sample_input.dataframe()
        # 输出CSV文件名和表名一致
        with output.filesystem().open(f"{table_name}.csv", "w") as stream:
            csv_writer = csv.writer(stream)
            csv_writer.writerow(df.schema.names)
            csv_writer.writerows(df.collect())
    # 注册任务到全局上下文
    globals()[f"export_{table_name}"] = export_task

注意事项

  • 代码完全符合Foundry代码仓库的开发规范,每个表的处理都是独立任务,支持并行调度,单表报错不会影响其他表的处理
  • 如果单表数据量极大,collect()导出单CSV可能触发内存超限,可以改用write_dataframe方法直接导出csv格式,会自动生成拆分的csv文件,稳定性更高
  • 所有输出路径、文件名自动根据表名生成,不会出现重名覆盖问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 03:06:03