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

Scala+Spark Streaming如何处理数组对象提取字段并计算时长?

在Scala + Spark Streaming中处理包含对象数组的列

实现思路

利用Spark内置的数组高阶函数transform遍历serviceDetails数组中的每个元素,对每个元素完成以下操作:

  1. 拼接stayBegin和stayEnd字段生成stayPeriod
  2. 将时间字符串转为秒级时间戳后计算差值,得到stayDuration
  3. 用struct构造新的对象结构,替换原数组中的元素

完整代码示例

1. 导入依赖与定义数据结构

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

object ServiceDetailsProcessing {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("ServiceDetailsProcessing")
      .master("local[*]") // 生产环境移除该行
      .getOrCreate()
    import spark.implicits._

    // 定义源数据Schema(从外部数据源读取时可直接复用)
    val serviceDetailSchema = StructType(Seq(
      StructField("serviceType", StringType),
      StructField("serviceOrder", LongType),
      StructField("stayOrder", LongType),
      StructField("stayBegin", StringType),
      StructField("stayEnd", StringType),
      StructField("locationID", StringType)
    ))
    val sourceSchema = StructType(Seq(
      StructField("serviceDetails", ArrayType(serviceDetailSchema))
    ))

2. 静态数据处理(测试验证用)

// 构造测试数据
    val testData = Seq(
      """{
        |  "serviceDetails": [
        |    {
        |      "serviceType": "xwFOisGAJbJlgpgodye",
        |      "serviceOrder": 20686918,
        |      "stayOrder": 14938272,
        |      "stayBegin": "2023-04-19T10:39:43",
        |      "stayEnd": "2023-04-19T11:39:43",
        |      "locationID": "NXPlsqagPcYMTPwJqErX"
        |    },
        |    {
        |      "serviceType": "wQmJTXOhzBAwbaatftsZ",
        |      "serviceOrder": 2949213,
        |      "stayOrder": 11157169,
        |      "stayBegin": "2023-04-19T10:39:43",
        |      "stayEnd": "2023-04-19T11:39:43",
        |      "locationID": "cJxXElbuuRVNMERFykpO"
        |    }
        |  ]
        |}""".stripMargin
    ).toDF("json")
      .select(from_json(col("json"), sourceSchema).alias("data"))
      .select("data.*")

    // 执行核心转换逻辑
    val resultDF = testData.withColumn(
      "serviceDetails",
      transform(
        col("serviceDetails"),
        elem => struct(
          concat(elem("stayBegin"), lit(" - "), elem("stayEnd")).alias("stayPeriod"),
          (unix_timestamp(elem("stayEnd")) - unix_timestamp(elem("stayBegin"))).alias("stayDuration")
        )
      )
    )

    // 查看转换结果
    resultDF.show(truncate = false)

3. Spark Streaming(Structured Streaming)场景适配

// 从Kafka读取流数据(示例,可替换为其他流数据源)
    val streamingDF = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "your-kafka-host:9092")
      .option("subscribe", "your-topic")
      .load()
      .select(from_json(col("value").cast(StringType), sourceSchema).alias("data"))
      .select("data.*")

    // 对流数据执行相同转换逻辑
    val processedStreamingDF = streamingDF.withColumn(
      "serviceDetails",
      transform(
        col("serviceDetails"),
        elem => struct(
          concat(elem("stayBegin"), lit(" - "), elem("stayEnd")).alias("stayPeriod"),
          (unix_timestamp(elem("stayEnd")) - unix_timestamp(elem("stayBegin"))).alias("stayDuration")
        )
      )
    )

    // 输出结果到控制台(生产环境可替换为HDFS、Kafka等存储)
    processedStreamingDF.writeStream
      .outputMode("append")
      .format("console")
      .option("truncate", "false")
      .start()
      .awaitTermination()
  }
}

注意事项

  • 如果时间字符串格式不是Spark默认支持的yyyy-MM-dd'T'HH:mm:ss,需要在unix_timestamp中显式指定格式,例如:
    unix_timestamp(elem("stayEnd"), "yyyy-MM-dd HH:mm:ss")
    
  • 若使用Spark版本低于2.4,transform函数不可用,此时需要自定义UDF来遍历处理数组元素。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 11:22:54