如何将温度曲线数据集转换为KMean等聚类算法适配格式
温度曲线数据集聚类预处理方案
针对你的大型温度曲线数据集,要适配KMeans等聚类算法,核心是将每个id对应的时间序列数据转换为固定长度的数值型特征向量,以下是具体步骤(以PySpark为例,适配大数据场景):
1. 解析字符串格式的temp列
你的temp列是字符串格式的嵌套列表,首先要把它解析成Spark能识别的数值型嵌套数组:
- 先定义temp列的嵌套结构Schema:
from pyspark.sql.types import ArrayType, StructType, StructField, DoubleType # 单个数据点的Schema:[毫秒数, 锅温, 环境温, 物体温] point_schema = StructType([ StructField("timestamp_ms", DoubleType(), nullable=True), StructField("pot_temp", DoubleType(), nullable=True), StructField("ambient_temp", DoubleType(), nullable=True), StructField("object_temp", DoubleType(), nullable=True) ]) # temp列的整体Schema:数组类型,每个元素是单个数据点 temp_schema = ArrayType(point_schema)
- 用
from_json解析字符串列:
from pyspark.sql.functions import from_json # 解析temp字符串为嵌套数组 df_parsed = df.withColumn("temp_parsed", from_json(df.temp, temp_schema)).drop("temp")
这一步会把原来的字符串转换成包含数值型字段的数组,解决你之前“数组元素是字符串”的问题。
2. 时间序列对齐与特征提取
因为每个id对应的时间序列长度可能不一致(有的数据点多,有的少),而KMeans要求每个样本的特征向量长度固定,所以需要做以下二选一处理:
方案A:固定时间间隔采样
如果时间序列的采样频率不统一,先按固定时间间隔(比如每1000ms)插值补全,得到长度一致的序列:
from pyspark.sql.functions import udf, explode, collect_list, struct import pandas as pd import numpy as np # 定义UDF:对单个id的时间序列按固定间隔插值 def interpolate_timeseries(points): # 转换为pandas DataFrame df = pd.DataFrame(points) if df.empty: return [] # 按时间排序 df = df.sort_values("timestamp_ms") # 生成固定间隔的时间轴(比如从0到最大时间,步长1000ms) max_ts = df["timestamp_ms"].max() new_ts = np.arange(0, max_ts + 1000, 1000) # 对每个温度维度插值 pot_interp = np.interp(new_ts, df["timestamp_ms"], df["pot_temp"]) ambient_interp = np.interp(new_ts, df["timestamp_ms"], df["ambient_temp"]) object_interp = np.interp(new_ts, df["timestamp_ms"], df["object_temp"]) # 合并成数组返回 return [list(row) for row in zip(new_ts, pot_interp, ambient_interp, object_interp)] interpolate_udf = udf(interpolate_timeseries, temp_schema) # 对每个id的时间序列插值 df_interpolated = df_parsed.groupBy("id").agg(collect_list("temp_parsed").alias("temp_list")) df_interpolated = df_interpolated.withColumn("temp_interpolated", interpolate_udf(df_interpolated.temp_list)).drop("temp_list")
之后可以把插值后的序列展开成一维向量,比如把所有锅温、环境温、物体温的插值结果拼接成一个长数组。
方案B:提取统计特征(更高效,适合大数据)
如果不需要保留完整序列趋势,直接对每个温度维度提取统计特征,生成固定长度的向量:
from pyspark.sql.functions import col, avg, max, min, stddev, count # 先把嵌套数组展开成行 df_exploded = df_parsed.select("id", explode(col("temp_parsed")).alias("point")) df_exploded = df_exploded.select( "id", col("point.timestamp_ms").alias("timestamp_ms"), col("point.pot_temp").alias("pot_temp"), col("point.ambient_temp").alias("ambient_temp"), col("point.object_temp").alias("object_temp") ) # 按id分组,计算每个维度的统计特征 df_features = df_exploded.groupBy("id").agg( avg("pot_temp").alias("pot_avg"), max("pot_temp").alias("pot_max"), min("pot_temp").alias("pot_min"), stddev("pot_temp").alias("pot_std"), avg("ambient_temp").alias("ambient_avg"), max("ambient_temp").alias("ambient_max"), min("ambient_temp").alias("ambient_min"), stddev("ambient_temp").alias("ambient_std"), avg("object_temp").alias("object_avg"), max("object_temp").alias("object_max"), min("object_temp").alias("object_min"), stddev("object_temp").alias("object_std"), count("timestamp_ms").alias("point_count") )
这样每个id对应一行,包含13个固定长度的特征,完全符合KMeans的输入要求。
3. 生成聚类所需的特征向量
用VectorAssembler把所有特征列合并成一个向量列:
from pyspark.ml.feature import VectorAssembler # 列出所有特征列名 feature_cols = [ "pot_avg", "pot_max", "pot_min", "pot_std", "ambient_avg", "ambient_max", "ambient_min", "ambient_std", "object_avg", "object_max", "object_min", "object_std", "point_count" ] assembler = VectorAssembler(inputCols=feature_cols, outputCol="features") df_final = assembler.transform(df_features)
最终的df_final中的features列就是KMeans可以直接输入的向量格式。
内容的提问来源于stack exchange,提问作者Adin
相关产品推荐
相关产品推荐

