I am trying to read a multiline file using pyspark in databricks and then flatten the file to get respective column which will invariably be stored as a delta table. I have successfully done it for 40 odd files but having performance issue with the last one.
This flatten part is taking 10+ mins to generate but the write part is running for ever.
#Flatten the dataframe to get all keys
def flatten(df):
complex_fields = dict([
(field.name, field.dataType)
for field in df.schema.fields
if isinstance(field.dataType, T.ArrayType) or isinstance(field.dataType, T.StructType)
])
qualify = list(complex_fields.keys())[0] + "_"
while len(complex_fields) != 0:
col_name = list(complex_fields.keys())[0]
if isinstance(complex_fields[col_name], T.StructType):
expanded = [F.col(col_name + '.' + k).alias(col_name + '_' + k)
for k in [ n.name for n in complex_fields[col_name]]
]
df = df.select("*", *expanded).drop(col_name)
elif isinstance(complex_fields[col_name], T.ArrayType):
df=df.withColumn(col_name,F.explode_outer(col_name))
complex_fields = dict([
(field.name, field.dataType)
for field in df.schema.fields
if isinstance(field.dataType, T.ArrayType) or isinstance(field.dataType, T.StructType)
])
for df_col_name in df.columns:
df = df.withColumnRenamed(df_col_name, df_col_name.replace(qualify, ""))
return df
df1=flatten(df)
To increase performance i have increased the cores and executor memory , not sure what else can be done. This is 1 single file but have multiple arrays, looking at the DAG looks like one input record is going to create millions of records in the target file.
Need help if anyone has faced similar issues.
I have attached the sql interpeter that is generated automatically. From stdout log all i can understand is Garbace Collector(GC) is trying to free up space , and also has multiple Full GC running.