Databricks中PySpark Streaming多传感器流数据聚合生成新传感器需求
传感器流数据处理需求
我每分钟将多个传感器的流数据接入Databricks平台。需要在每次PySpark Streaming加载任务中,若传感器'ABC'和'DEF'同时存在,则创建新传感器'PQRS',其数值为ABC和DEF数值的平均值。
输入数据(1分钟流)
| sensor_name | value | timestamp |
|---|---|---|
| ABC | 10 | 2023-11-02T11:49:32.028Z |
| DEF | 20 | 2023-11-02T11:49:32.028Z |
| GHI | 12 | 2023-11-02T11:49:32.028Z |
输出数据
| sensor_name | value | timestamp |
|---|---|---|
| ABC | 10 | 2023-11-02T11:49:32.028Z |
| DEF | 20 | 2023-11-02T11:49:32.028Z |
| PQRS | 15 | 2023-11-02T11:49:32.028Z |
| GHI | 12 | 2023-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()
逻辑说明
- 分组聚合:按
timestamp对每个批次的流数据分组,收集所有传感器条目,并通过计数判断当前批次是否同时包含ABC和DEF。 - 扩展传感器列表:使用UDF遍历传感器列表,提取ABC和DEF的数值,计算平均值后添加PQRS条目(仅当两个传感器都存在时)。
- 展开数据:将扩展后的传感器列表展开为原始的行结构,保证输出格式与输入一致。
内容的提问来源于stack exchange,提问作者paul saju
相关产品推荐
相关产品推荐

