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

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"))

问题分析与解决方案

核心错误点

  1. matching_rules的Schema定义错误:目标JSON中matching_rules是数组类型(用[]包裹),但你错误地将其定义为StructType(Array(...)),正确类型应为ArrayType(StructType(...)),表示“结构体数组”。
  2. 数组字段提取方式错误:即便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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 06:25:27