如何对AWS Glue的Dynamic Frame按name、key排序?解析drop函数作用
AWS Glue Job 排序实现与代码解释
一、实现按"name"、"key"排序的方案
场景1:全局排序(整个数据集按规则排序)
如果是要对整个Dynamic Frame数据先按name列排序,name相同时按key列排序,需要先把Dynamic Frame转为Spark DataFrame处理,排序后再转回Dynamic Frame:
from awsglue.dynamicframe import DynamicFrame # 假设你的Dynamic Frame实例是dyf spark_df = dyf.toDF() # 按name升序、key升序排序;要降序的话用f.desc("name")、f.desc("key") sorted_spark_df = spark_df.orderBy(f.col("name"), f.col("key")) # 转回Dynamic Frame sorted_dyf = DynamicFrame.fromDF(sorted_spark_df, glue_context, "sorted_dynamic_frame")
场景2:分组后组内按规则排序并取第一条
如果是要按某些字段分组,组内先按name排序、name相同按key排序,然后保留每组的第一条记录,可以基于你提供的函数修改:
def get_top_sorted_records(data_frame, group_keys): original_columns = data_frame.columns # 定义窗口:按group_keys分组,组内先按name排序,再按key排序 window_spec = w.partitionBy(*group_keys).orderBy(f.col("name"), f.col("key")) output_df = data_frame.withColumn("row_num", f.row_number().over(window_spec)). \ filter(f.col("row_num") == 1). \ drop(f.col("row_num")). \ select(original_columns) return output_df # 使用示例:按"category"分组,处理Dynamic Frame spark_df = dyf.toDF() result_spark_df = get_top_sorted_records(spark_df, ["category"]) result_dyf = DynamicFrame.fromDF(result_spark_df, glue_context, "result_dynamic_frame")
二、drop函数的作用
原代码里的drop(f.col("row_num"))是用来删除临时生成的row_num列:
- 为了给分组内的记录编号,我们用
withColumn新增了row_num列,存储每条记录在组内的行号; - 筛选出
row_num == 1的目标记录后,这个临时列就没用了,用drop删掉它,能保证最终输出的DataFrame列结构和原始数据一致(后续的select(original_columns)也进一步确认了这一点)。
内容的提问来源于stack exchange,提问作者light
相关产品推荐
相关产品推荐

