PySpark中Struct类型无法使用explode函数的问题求助
Got it, let's break down why you're hitting that error and how to fix it.
Why the Error Happens
Looking at your JSON schema, the values column is a struct type (it contains nested fields like Contrex_table, UFI_Table, etc.), not an array or map. Spark's explode() function only works on array or map columns—this is exactly why you're seeing the "data type mismatch" message.
Solution: Convert Struct to Map Then Explode
To unpack the values struct into usable rows, we first convert the struct to a map (where keys are field names and values are the corresponding arrays), then explode that map. After that, we can explode the nested arrays to get individual records tied to their original identifier.
Here's the revised, working code:
from pyspark.sql.functions import explode, create_map, lit from itertools import chain # Your existing code up to df3 start_df = spark.read.json(f'/mnt/bronze/products/**/*.json') df = start_df.select('_embedded.*') df1 = df.select(explode(df.items).alias('required')) df2 = df1.select('required.*') df3= df2.select("identifier","values") df3.display() # Step 1: Dynamically get all field names from the 'values' struct values_field_names = df3.select("values.*").columns # Step 2: Create a map where each entry is (field_name, field_value) map_expression = create_map( list(chain.from_iterable([lit(field), df3[f"values.{field}"]] for field in values_field_names)) ) # Step 3: Explode the map to turn each struct field into a separate row df_4 = df3.select("identifier", explode(map_expression).alias("field_name", "field_values")) # Step 4: Explode the nested array in 'field_values' to get individual items df_5 = df_4.select("identifier", "field_name", explode("field_values").alias("value_details")) # Step 5: Expand the 'value_details' struct to expose all nested columns df_final = df_5.select("identifier", "field_name", "value_details.*") df_final.display()
What Each Step Does
- Step 1: We dynamically fetch all field names from the
valuesstruct, so you don't have to hardcode them if the schema changes later. - Step 2: Convert the struct into a map, pairing each field name with its value (which is an array in your schema).
- Step 3: Explode the map—this creates a row for every field in the
valuesstruct, keeping the originalidentifierlinked to each entry. - Step 4: Explode the array in
field_valuesto unpack each individual item in the array into its own row. - Step 5: Finally, expand the
value_detailsstruct to reveal all its nested columns likedata,locale, andscope.
Alternative: Explicit Unpivoting with stack
If your values struct has a fixed set of fields and you prefer a more explicit approach, use the stack function to unpivot the struct:
from pyspark.sql.functions import explode, stack # Reuse values_field_names from earlier df_4 = df3.select( "identifier", stack( len(values_field_names), *[lit(field), df3[f"values.{field}"] for field in values_field_names] ).alias("field_name", "field_values") ) # Proceed with exploding and expanding as before df_5 = df_4.select("identifier", "field_name", explode("field_values").alias("value_details")) df_final = df_5.select("identifier", "field_name", "value_details.*") df_final.display()
This achieves the same result as the map approach but uses stack for a more direct, explicit unpivot.
内容的提问来源于stack exchange,提问作者Greencolor

