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
相关产品推荐
相关产品推荐

