如何在AWS Glue Job中按学生ID提取最新记录并输出
在AWS Glue Job中提取每个学生的最新记录(按SQNo最大值)
我之前处理学生数据时遇到过完全一样的需求,给你分享两种靠谱的实现方式,不管是用Spark DataFrame还是直接操作DynamicFrame都能搞定,应该能解决你之前转DataFrame没成功的问题~
方法一:用Spark DataFrame + 窗口函数(推荐)
这种方法逻辑清晰,处理大数据量也高效,步骤如下:
读取CSV文件到DynamicFrame
首先用Glue的API读取你的CSV数据,注意根据实际情况调整分隔符、是否有表头这些参数:from awsglue.context import GlueContext from pyspark.context import SparkContext sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) # 读取CSV,这里假设你的数据是空格分隔,第一行是表头 dynamic_frame = glueContext.create_dynamic_frame.from_options( connection_type="s3", connection_options={"paths": ["s3://your-input-bucket/path/to/csv/"]}, format="csv", format_options={ "separator": " ", "withHeader": True, "quoteChar": '"' } )转换为Spark DataFrame
这一步很简单,直接调用toDF()方法就行:df = dynamic_frame.toDF() # 可以先打印下数据结构确认:df.printSchema()用窗口函数筛选每个StudentID的最新记录
我们通过窗口函数给每个StudentID分组内的记录按SQNo降序排,然后取每组的第一条:from pyspark.sql.window import Window import pyspark.sql.functions as F # 定义窗口:按StudentID分组,按SQNo降序排序 window_spec = Window.partitionBy("StudentID").orderBy(F.desc("SQNo")) # 给每条记录加行号,每组内第一条(SQNo最大)的行号是1 df_with_row_num = df.withColumn("row_num", F.row_number().over(window_spec)) # 筛选行号为1的记录,就是每个学生的最新记录 latest_records_df = df_with_row_num.filter(F.col("row_num") == 1).drop("row_num")输出为CSV或数据表
如果你想转回DynamicFrame输出(比如用Glue的内置写操作):from awsglue.dynamicframe import DynamicFrame latest_records_dynamic = DynamicFrame.fromDF(latest_records_df, glueContext, "latest_records") # 写入S3为CSV glueContext.write_dynamic_frame.from_options( frame=latest_records_dynamic, connection_type="s3", connection_options={"path": "s3://your-output-bucket/path/to/output/"}, format="csv", format_options={"separator": " ", "withHeader": True} ) # 或者写入数据湖表(比如Glue Data Catalog中的表) # glueContext.write_dynamic_frame.from_catalog( # frame=latest_records_dynamic, # database="your-db", # table_name="your-output-table" # )
方法二:分组取最大值 + 关联原表
如果不想用窗口函数,也可以先分组获取每个StudentID的最大SQNo,再和原表关联得到完整记录:
# 第一步读取和转DataFrame和上面一样,直接从第二步开始 max_sqno_df = df.groupBy("StudentID").agg(F.max("SQNo").alias("max_SQNo")) # 关联原表,匹配StudentID和对应的最大SQNo latest_records_df = df.join( max_sqno_df, (df.StudentID == max_sqno_df.StudentID) & (df.SQNo == max_sqno_df.max_SQNo), "inner" ).drop(max_sqno_df.StudentID) # 去掉重复的StudentID列
为什么你之前转DataFrame可能没成功?
大概率是这两个原因:
- 窗口函数的
orderBy方向错了:如果用了升序,取到的会是SQNo最小的记录,而不是最大的 - 没有正确处理分组后的关联:如果只分组取max(SQNo),会丢失其他字段,必须再和原表关联才能拿到完整记录
内容的提问来源于stack exchange,提问作者Amit
相关产品推荐
相关产品推荐

