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

如何在AWS Glue中按SrlNo顺序拆分DataFrame为多个CSV文件

按指定记录数有序拆分DataFrame为多个CSV文件(Glue 3.0/Spark 3.1)

问题背景

现有如下DataFrame:

+--------+------+
|Name    | SrlNo|
+--------+------+
|Sweden  | 1    |
|Albania | 2    |
|India   | 3    |
|Iceland | 4    |
|Finland | 5    |
|Denmark | 6    |
|Algeria | 8    |
|Andorra | 9    |
|Norway  | 10   |
+-------+-------|

需要按指定记录数拆分并保存为多个CSV文件,但当前使用的Glue代码生成的文件数据乱序。

当前代码

finalCount=dynamicFrame.count()
records_per_file=14701
partition_count = math.ceil(finalCount / records_per_file)
if partition_count < 1:
    partition_count = 1

dynamicFrame = dynamicFrame.repartition(partition_count)
glueContext.write_dynamic_frame.from_options(
    frame=dynamicFrame,
    connection_type="s3",
    connection_options={
        "path": "S3_Path",
        'groupFiles': 'inPartition', 'groupSize': '10485760'
    },
    format="csv",
    format_options={
        "optimizePerformance": True, 
        "separator": ","
        },
    transformation_ctx="AmazonS3_",
)

当前错误输出

CSV 1:

+--------+------+
|Name    | SrlNo|
+--------+------+
|Sweden  | 1    |
|India   | 3    |
|Finland | 5    |
|Denmark | 6    |
|Andorra | 9    |
+-------+-------|

CSV 2:

+--------+------+
|Name    | SrlNo|
+--------+------+
|Albania | 2    |
|India   | 3    |
|Iceland | 4    |
|Algeria | 8    |
|Norway  | 10   |
+-------+-------|

期望输出

CSV 1:

+--------+------+
|Name    | SrlNo|
+--------+------+
|Sweden  | 1    |
|Albania | 2    |
|India   | 3    |
|Iceland | 4    |
+-------+-------|

CSV 2:

+--------+------+
|Name    | SrlNo|
+--------+------+
|Finland | 5    |
|Denmark | 6    |
|Algeria | 8    |
|Andorra | 9    |
|Norway  | 10   |
+-------+-------|

解决方案

问题根源在于repartition是随机分配数据到分区的,不会保留原有顺序,且写入时未保证分区内数据有序。以下是修正步骤:

1. 转换为Spark DataFrame

Glue DynamicFrame的排序和分区控制不如Spark DataFrame灵活,先进行格式转换:

from pyspark.sql import Window
import pyspark.sql.functions as F
import math
from awsglue.dynamicframe import DynamicFrame

# 将DynamicFrame转为Spark DataFrame
df = dynamicFrame.toDF()

2. 全局排序并计算目标分区

按SrlNo排序后,用窗口函数给每行分配行号,再根据records_per_file计算该行所属的分区:

records_per_file = 14701
finalCount = df.count()
partition_count = math.ceil(finalCount / records_per_file)
if partition_count < 1:
    partition_count = 1

# 按SrlNo全局排序,生成连续行号
window_spec = Window.orderBy("SrlNo")
df = df.withColumn("row_num", F.row_number().over(window_spec))

# 计算每行所属的目标分区ID(从0开始)
df = df.withColumn("partition_id", F.floor((F.col("row_num") - 1) / records_per_file))

3. 按分区ID重新分区并写入CSV

按partition_id分区,确保同一分区内的行是连续有序的;写入时关闭性能优化选项以保留顺序:

# 按partition_id分区,保证每个分区对应一个有序的CSV文件段
df_partitioned = df.repartition(partition_count, "partition_id")

# 转换回DynamicFrame(也可直接用Spark原生write方法)
dynamicFrame_partitioned = DynamicFrame.fromDF(df_partitioned, glueContext, "dynamicFrame_partitioned")

# 写入CSV,关键设置:关闭optimizePerformance以保留排序结果
glueContext.write_dynamic_frame.from_options(
    frame=dynamicFrame_partitioned,
    connection_type="s3",
    connection_options={
        "path": "S3_Path",
        'groupFiles': 'inPartition',  # 每个分区生成一个独立文件
        'groupSize': '10485760'
    },
    format="csv",
    format_options={
        "optimizePerformance": False,  # 必须关闭,否则会打乱已排好的顺序
        "separator": ",",
        "header": True  # 保留表头
        },
    transformation_ctx="AmazonS3_",
)

关键说明

  • 全局排序:必须先按SrlNo完成全局排序,确保后续分区的记录是连续的有序段
  • 分区ID计算:通过行号与单文件记录数的比值确定分区,保证每个分区内的记录范围连续
  • 关闭性能优化:optimizePerformance会触发Spark的底层优化,打乱已排好的顺序,所以必须禁用
  • 单分区单文件:groupFiles: 'inPartition'确保每个分区生成一个独立CSV文件,避免多段数据合并

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 08:36:14