如何在AWS Glue中按日期对DynamicFrame正确排序?
解决AWS Glue中按日期字段排序CSV的问题
你说得没错,birthdate字段被识别为字符串类型,字符串排序是按字符顺序而非实际日期顺序,导致排序结果不符合预期。需要先将字符串转换为日期类型,再进行排序,以下是两种可行的解决方案:
方案1:读取CSV时直接指定Schema,解析日期类型
在读取阶段就明确字段类型,避免后续转换步骤,效率更高:
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from awsglue.dynamicframe import DynamicFrame from pyspark.sql.types import StructType, StructField, StringType, DateType sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) # 补充获取Job参数的步骤,原代码遗漏 args = getResolvedOptions(sys.argv, ['JOB_NAME']) job.init(args["JOB_NAME"], args) # 定义自定义Schema,明确birthdate为日期类型 custom_schema = StructType([ StructField("firstname", StringType(), True), StructField("lastname", StringType(), True), StructField("birthdate", DateType(), True) ]) myCSV = glueContext.create_dynamic_frame.from_options( format_options={ "quoteChar": '"', "withHeader": True, "separator": ",", "optimizePerformance": False, "dateFormat": "MM/dd/yyyy" # 指定匹配CSV的日期格式 }, connection_type="s3", format="csv", connection_options={ "paths": [ "S3-PATH-TO-CSV" ] }, transformation_ctx="myCSV", schema=custom_schema # 应用自定义Schema ) # 直接按日期类型的birthdate字段排序 sortedCSV = myCSV.toDF().sort("birthdate") # 重新分区为单个文件并写入S3 repartitionedDF = DynamicFrame.fromDF(sortedCSV, glueContext, "repartitionedDF").repartition(1) outputCSV = glueContext.write_dynamic_frame.from_options( frame=repartitionedDF, connection_type="s3", format="csv", connection_options={ "path": "OUTPATH-S3", "partitionKeys": [], }, format_options={"withHeader": True}, # 写入时保留表头 transformation_ctx="outputCSV", ) job.commit()
方案2:读取后将字符串转换为日期类型
如果无法在读取阶段指定Schema,可在读取完成后用Spark函数转换字段类型:
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from awsglue.dynamicframe import DynamicFrame from pyspark.sql.functions import to_date sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) args = getResolvedOptions(sys.argv, ['JOB_NAME']) job.init(args["JOB_NAME"], args) myCSV = glueContext.create_dynamic_frame.from_options( format_options={ "quoteChar": '"', "withHeader": True, "separator": ",", "optimizePerformance": False, }, connection_type="s3", format="csv", connection_options={ "paths": [ "S3-PATH-TO-CSV" ] }, transformation_ctx="myCSV", ) # 将字符串类型的birthdate转换为日期类型 df = myCSV.toDF() df_with_date = df.withColumn("birthdate", to_date(df["birthdate"], "MM/dd/yyyy")) # 按日期字段排序 sortedCSV = df_with_date.sort("birthdate") # 重新分区并写入S3 repartitionedDF = DynamicFrame.fromDF(sortedCSV, glueContext, "repartitionedDF").repartition(1) outputCSV = glueContext.write_dynamic_frame.from_options( frame=repartitionedDF, connection_type="s3", format="csv", connection_options={ "path": "OUTPATH-S3", "partitionKeys": [], }, format_options={"withHeader": True}, transformation_ctx="outputCSV", ) job.commit()
关键注意事项
- 日期格式字符串
"MM/dd/yyyy"会自动兼容单数字的月份/日期(如9/15/1994),无需修改为M/d/yyyy。 - 写入CSV时必须添加
format_options={"withHeader": True},否则输出文件会丢失表头。 - 原代码遗漏了
args = getResolvedOptions(sys.argv, ['JOB_NAME']),这是获取Glue Job参数的必要步骤,需补充。
内容的提问来源于stack exchange,提问作者uzluisf
相关产品推荐
相关产品推荐

