Scala+Spark Streaming如何处理数组对象提取字段并计算时长?
在Scala + Spark Streaming中处理包含对象数组的列
实现思路
利用Spark内置的数组高阶函数transform遍历serviceDetails数组中的每个元素,对每个元素完成以下操作:
- 拼接
stayBegin和stayEnd字段生成stayPeriod - 将时间字符串转为秒级时间戳后计算差值,得到
stayDuration - 用
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
相关产品推荐
相关产品推荐

