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

PySpark读取Kafka JSON数据时数组字段引用歧义问题如何解决

问题根源

你收到的歧义错误和数组类型无关,是你自定义的JSON Schema结构错误导致的:你的Schema里value结构体下出现了两个重名的entity字段,一个仅携带type属性,另一个仅携带sources属性,Spark解析到value.entity时无法判断你要访问哪个结构体,因此抛出异常。

第一步:修正Schema

正确的Schema应该是entity一个结构体下同时包含type和sources两个子字段,参考定义如下:

from pyspark.sql.types import *

# 最内层的items元素结构
item_struct = StructType([
    StructField("identifier", StringType(), True)
])
# sources元素结构
source_struct = StructType([
    StructField("items", ArrayType(item_struct), True)
])
# entity结构,同时包含type和sources
entity_struct = StructType([
    StructField("type", StringType(), True),
    StructField("sources", ArrayType(source_struct), True)
])
# 最外层的value结构
schema = StructType([
    StructField("action", StringType(), True),
    StructField("id", StringType(), True),
    StructField("epoch", LongType(), True),
    StructField("entity", entity_struct, True)
])

用这个Schema重新解析JSON,就不会再出现字段歧义问题。

第二步:提取嵌套数组中的字段

修正Schema后,针对sources、items两层数组结构,你可以根据需求选择两种处理方式:

方式1:保留数组格式,直接提取所有identifier到数组中

不需要展开数组,直接用高阶函数提取所有identifier:

from pyspark.sql.functions import from_json, col, transform, flatten

value_df = df.select(
    from_json(col("value").cast("string"), schema).alias("value"), 
)

result_df = value_df.select(
    "value.id", 
    "value.action", 
    "value.epoch",
    "value.entity.type",
    # 先提取所有items里的identifier,再把两层数组拍平成一维数组
    flatten(transform(col("value.entity.sources"), 
                      lambda s: transform(s.items, lambda i: i.identifier))
           ).alias("identifiers")
)

此时identifiers字段是字符串数组,包含该条消息里所有的identifier值。

方式2:展开数组,每个identifier对应一行

如果需要把每个identifier拆成独立的行,用explode函数逐层展开数组即可:

from pyspark.sql.functions import from_json, col, explode

value_df = df.select(
    from_json(col("value").cast("string"), schema).alias("value"), 
)

# 先展开sources数组
explode_source = value_df.select(
    "value.id", 
    "value.action", 
    "value.epoch",
    "value.entity.type",
    explode("value.entity.sources").alias("source")
)
# 再展开items数组
explode_item = explode_source.select(
    "id", "action", "epoch", "type",
    explode("source.items").alias("item")
)
# 最后提取identifier
result_df = explode_item.select(
    "id", "action", "epoch", "type",
    col("item.identifier").alias("identifier")
)

内容的提问来源于stack exchange,提问作者Freid001

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 20:27:04