Performance of PySpark DataFrames vs Glue DynamicFrames

Viewed 241

So I recently started using Glue and PySpark for the first time. The task was to create a Glue job that does the following:

  1. Load data from parquet files residing in an S3 bucket
  2. Apply a filter to the data
  3. Add a column, the value of which is derived from 2 other columns
  4. Write the result to S3

Since the data is going from S3 to S3, I assumed that Glue DynamicFrames should be a decent fit for this, and I came up with the following code:

def AddColumn(r):
   if r["option_type"] == 'S': 
       r["option_code_derived"]= 'S'+ r["option_code_4"]
   elif r["option_type"] == 'P': 
       r["option_code_derived"]= 'F'+ r["option_code_4"][1:]
   elif r["option_type"] == 'L':
       r["option_code_derived"]= 'P'+ r["option_code_4"]
   else:  
       r["option_code_derived"]= None
       
   return r

glueContext = GlueContext(create_spark_context(role_arn=args['role_arn']))
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

inputGDF = glueContext.create_dynamic_frame_from_options(connection_type = "s3", connection_options = {"paths": [source_path], "recurse" : True}, format = source_format, additional_options = {"useS3ListImplementation":True})

filtered_gdf = Filter.apply(frame = inputGDF, f = lambda x: x["my_filter_column"] in ['50','80'])

additional_column_gdf = Map.apply(frame = filtered_gdf, f = AddColumn) 

gdf_mapped = ApplyMapping.apply(frame = additional_column_gdf, mappings = mappings, transformation_ctx = "gdf_mapped") 

glueContext.purge_s3_path(full_target_path_purge, {"retentionPeriod": 0})

outputGDF = glueContext.write_dynamic_frame.from_options(frame = gdf_mapped, connection_type = "s3", connection_options = {"path": full_target_path}, format = target_format)

This works but takes a very long time (just short of 10 hours with 20 G1.X workers). Now, the dataset is quite large (almost 2 billion records, over 400 GB), but this was still unexpected (to me at least).

Then I gave it another try, this time with PySpark DataFrames instead of DynamicFrames. The code looks like the following:

glueContext = GlueContext(create_spark_context(role_arn=args['role_arn'], source_bucket=args['s3_source_bucket'], target_bucket=args['s3_target_bucket']))
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

df = spark.read.parquet(full_source_path)

df_filtered = df.filter( (df.model_key_status  == '50') | (df.model_key_status  == '80') )

df_derived = df_filtered.withColumn('option_code_derived', 
     when(df_filtered.option_type == "S", concat(lit('S'), df_filtered.option_code_4))
    .when(df_filtered.option_type == "P", concat(lit('F'), df_filtered.option_code_4[2:42]))
    .when(df_filtered.option_type == "L", concat(lit('P'), df_filtered.option_code_4))
    .otherwise(None))

glueContext.purge_s3_path(full_purge_path, {"retentionPeriod": 0})

df_reorderered = df_derived.select(target_columns)

df_reorderered.write.parquet(full_target_path, mode="overwrite")

This also works, but with otherwise identical settings (20 workers of type G1.X, same dataset), this takes less than 20 minutes.

My question is: Where does this massive difference in performance between DynamicFrames and DataFrames come from? Was I doing something fundamentally wrong in the first try?

0 Answers
Related