You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark中Struct类型无法使用explode函数的问题求助

Fixing the "explode requires array/map type" Error in Spark JSON Processing

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 values struct, 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 values struct, keeping the original identifier linked to each entry.
  • Step 4: Explode the array in field_values to unpack each individual item in the array into its own row.
  • Step 5: Finally, expand the value_details struct to reveal all its nested columns like data, locale, and scope.

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 16:40:50