Spark Structured Streaming解析Kafka JSON时tag字段值为null的问题求助
问题:Spark Structured Streaming提取Kafka消息中
matching_rules.tag字段全为null 目标JSON结构
{ "data": { "created_at": "***", "id": "***", "text": "***" }, "matching_rules": [ { "id": "***", "tag": "***" } ] }
当前错误的Schema定义
val DFschema = StructType(Array( StructField("data", StructType(Array( StructField("created_at", TimestampType), StructField("text", StringType)))), StructField("matching_rules", StructType(Array( StructField("tag", StringType) ))) ))
当前处理代码
val kafkaDF: DataFrame = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", servers) .option("failOnDataLoss", "false") .option("subscribe", topics) .option("startingOffsets", "earliest") .load() .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") .select(col("key"), from_json($"value", DFschema).alias("structdata")) .select($"key", $"structdata.data".getField("created_at").alias("created_at"), $"structdata.data".getField("text").alias("text"), $"structdata.matching_rules".getField("tag").alias("topic") ) .withColumn("hour", date_format(col("created_at"), "HH")) .withColumn("date", date_format(col("created_at"), "yyyy-MM-dd"))
问题分析与解决方案
核心错误点
matching_rules的Schema定义错误:目标JSON中matching_rules是数组类型(用[]包裹),但你错误地将其定义为StructType(Array(...)),正确类型应为ArrayType(StructType(...)),表示“结构体数组”。- 数组字段提取方式错误:即便Schema修正,直接对数组列调用
getField("tag")也无法获取值,需先定位数组元素(如取第一个)或展开数组后再提取字段。
修正后的Schema
val DFschema = StructType(Array( // data是单个结构体,用StructType包裹字段 StructField("data", StructType(Seq( StructField("created_at", TimestampType), StructField("text", StringType) ))), // matching_rules是结构体数组,用ArrayType包裹StructType StructField("matching_rules", ArrayType(StructType(Seq( StructField("tag", StringType), StructField("id", StringType) // 可选,不需要可删除 )))) ))
修正后的提取逻辑
根据matching_rules数组的元素数量,有两种处理方式:
方式1:数组仅含单个元素(取第一个元素的tag)
如果每个消息的matching_rules固定只有一个元素,用element_at函数提取第一个元素的tag:
val kafkaDF: DataFrame = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", servers) .option("failOnDataLoss", "false") .option("subscribe", topics) .option("startingOffsets", "earliest") .load() .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") .select(col("key"), from_json($"value", DFschema).alias("structdata")) .select($"key", $"structdata.data.created_at".alias("created_at"), // 简化嵌套字段访问 $"structdata.data.text".alias("text"), element_at($"structdata.matching_rules", 1).getField("tag").alias("topic") // 取数组第一个元素的tag ) .withColumn("hour", date_format(col("created_at"), "HH")) .withColumn("date", date_format(col("created_at"), "yyyy-MM-dd"))
方式2:数组含多个元素(展开数组后提取)
如果matching_rules可能有多个元素,用explode将数组展开为多行,每行对应一个数组元素:
val kafkaDF: DataFrame = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", servers) .option("failOnDataLoss", "false") .option("subscribe", topics) .option("startingOffsets", "earliest") .load() .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") .select(col("key"), from_json($"value", DFschema).alias("structdata")) .select($"key", $"structdata.data.created_at".alias("created_at"), $"structdata.data.text".alias("text"), explode($"structdata.matching_rules").alias("rule") // 展开数组 ) .select($"key", $"created_at", $"text", $"rule.tag".alias("topic")) // 提取展开后的tag .withColumn("hour", date_format(col("created_at"), "HH")) .withColumn("date", date_format(col("created_at"), "yyyy-MM-dd"))
内容的提问来源于stack exchange,提问作者parmeni4
相关产品推荐
相关产品推荐

