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

