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

Databricks中PySpark Streaming多传感器流数据聚合生成新传感器需求

传感器流数据处理需求

我每分钟将多个传感器的流数据接入Databricks平台。需要在每次PySpark Streaming加载任务中,若传感器'ABC'和'DEF'同时存在,则创建新传感器'PQRS',其数值为ABC和DEF数值的平均值。

输入数据(1分钟流)

sensor_namevaluetimestamp
ABC102023-11-02T11:49:32.028Z
DEF202023-11-02T11:49:32.028Z
GHI122023-11-02T11:49:32.028Z

输出数据

sensor_namevaluetimestamp
ABC102023-11-02T11:49:32.028Z
DEF202023-11-02T11:49:32.028Z
PQRS152023-11-02T11:49:32.028Z
GHI122023-11-02T11:49:32.028Z

解决方案代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum, when, lit, explode, struct, collect_list
from pyspark.sql.types import StructType, StringType, DoubleType, TimestampType, ArrayType
from pyspark.sql.functions import udf

# 初始化SparkSession(Databricks环境中可省略)
spark = SparkSession.builder.appName("SensorStreamPQRS").getOrCreate()

# 定义输入Schema
input_schema = StructType() \
    .add("sensor_name", StringType()) \
    .add("value", DoubleType()) \
    .add("timestamp", TimestampType())

# 读取流数据(示例为文件流,可替换为Kafka、Event Hubs等数据源)
stream_df = spark.readStream \
    .schema(input_schema) \
    .format("csv") \
    .option("header", "true") \
    .load("/databricks/driver/sensor_stream")

# 按时间戳分组,收集传感器列表并判断是否同时存在ABC和DEF
grouped_df = stream_df.groupBy("timestamp") \
    .agg(
        collect_list(struct("sensor_name", "value")).alias("sensor_entries"),
        when(
            (sum(when(col("sensor_name") == "ABC", 1).otherwise(0)) >= 1) &
            (sum(when(col("sensor_name") == "DEF", 1).otherwise(0)) >= 1),
            lit(True)
        ).otherwise(lit(False)).alias("has_both_sensors")
    )

# 自定义UDF:扩展传感器列表,添加PQRS(如果条件满足)
def add_pqrs_sensor(sensor_list, has_both):
    abc_val = None
    def_val = None
    # 提取ABC和DEF的值
    for entry in sensor_list:
        if entry.sensor_name == "ABC":
            abc_val = entry.value
        elif entry.sensor_name == "DEF":
            def_val = entry.value
    # 复制原列表
    extended_list = list(sensor_list)
    # 添加PQRS
    if has_both and abc_val is not None and def_val is not None:
        extended_list.append({"sensor_name": "PQRS", "value": (abc_val + def_val)/2})
    return extended_list

# 注册UDF
add_pqrs_udf = udf(add_pqrs_sensor, ArrayType(StructType().add("sensor_name", StringType()).add("value", DoubleType())))

# 生成最终结果
result_df = grouped_df.withColumn("extended_sensors", add_pqrs_udf(col("sensor_entries"), col("has_both_sensors"))) \
    .select("timestamp", explode(col("extended_sensors")).alias("sensor")) \
    .select("sensor.sensor_name", "sensor.value", "timestamp")

# 输出到控制台(生产环境可替换为Delta Lake等持久化存储)
query = result_df.writeStream \
    .outputMode("append") \
    .format("console") \
    .option("truncate", "false") \
    .start()

query.awaitTermination()

逻辑说明

  1. 分组聚合:按timestamp对每个批次的流数据分组,收集所有传感器条目,并通过计数判断当前批次是否同时包含ABC和DEF。
  2. 扩展传感器列表:使用UDF遍历传感器列表,提取ABC和DEF的数值,计算平均值后添加PQRS条目(仅当两个传感器都存在时)。
  3. 展开数据:将扩展后的传感器列表展开为原始的行结构,保证输出格式与输入一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 22:05:56