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

PySpark DataFrame坐标数据转换:按行程聚合生成轨迹

高效实现PySpark轨迹数据聚合

直接用Spark的分布式分组聚合API就能解决,完全不需要遍历每行,适合大数据量场景:

核心思路

按Id、Start、Stop这三个行程唯一标识字段分组,将每组内的Lat(纬度)和Lon(经度)组合成坐标对,再收集成数组作为轨迹字段。

具体代码实现

1. 导入所需函数

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, collect_list, struct

2. 创建示例DataFrame(模拟你的数据)

spark = SparkSession.builder.appName("TrajectoryAggregation").getOrCreate()

data = [
    (1, 40.5, 40, "A", "B"),
    (1, 41.0, 45, "A", "B"),
    (1, 40.5, 40, "A", "B"),
    (2, 31.4, 59, "A", "C"),
    (2, 34.5, 60, "A", "C"),
    (2, 37.0, 61, "A", "C")
]

df = spark.createDataFrame(data, ["Id", "Lat", "Lon", "Start", "Stop"])

3. 执行分组聚合生成轨迹

trajectory_df = df.groupBy("Id", "Start", "Stop") \
                  .agg(collect_list(struct(col("Lat"), col("Lon"))).alias("Trajectory"))

trajectory_df.show(truncate=False)

输出结果

+---+-----+----+------------------------------------+
|Id |Start|Stop|Trajectory                          |
+---+-----+----+------------------------------------+
|1  |A    |B   |[{40.5, 40}, {41.0, 45}, {40.5, 40}]|
|2  |A    |C   |[{31.4, 59}, {34.5, 60}, {37.0, 61}]|
+---+-----+----+------------------------------------+

关键说明

  • struct(col("Lat"), col("Lon")):把纬度和经度打包成一个结构体,方便后续作为单个坐标点存储
  • collect_list(...):将每组内的所有坐标点收集成一个数组,这个函数是分布式执行的,Spark会自动在各个节点并行处理分区内的数据,避免单节点遍历的性能瓶颈
  • 如果需要保证轨迹点的顺序(比如按时间先后),可以先对原DataFrame按Id和时间字段(如果有的话)排序,再执行聚合:
    # 假设存在时间字段timestamp,先排序再聚合
    sorted_df = df.orderBy("Id", "timestamp")
    trajectory_df = sorted_df.groupBy("Id", "Start", "Stop") \
                             .agg(collect_list(struct(col("Lat"), col("Lon"))).alias("Trajectory"))
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 10:01:42