如何在Polars中按分区执行滚动窗口计算?以航线方位计算为例
解决Polars中按飞机分区计算方位的方案及优化建议
核心思路
Polars中rolling()与over()的组合存在限制,但计算飞机方位仅需当前点与前一个时间点的GPS坐标,因此可以用shift(1).over("plane_no")获取每组内的前序坐标,再结合方位角公式完成计算。无需依赖rolling窗口,这种方式更高效且适配Polars的分组逻辑。
具体实现方案
1. 基础方案:用shift()+over()结合自定义函数
先确保数据按飞机编号和时间戳排序,再获取每组内的前序GPS坐标,最后调用方位计算函数:
import polars as pl import math # 定义初始方位角计算函数(从点A到点B的方位,范围0-360度) def calculate_bearing(lat1, lon1, lat2, lon2): lat1_rad = math.radians(lat1) lon1_rad = math.radians(lon1) lat2_rad = math.radians(lat2) lon2_rad = math.radians(lon2) delta_lon = lon2_rad - lon1_rad x = math.sin(delta_lon) * math.cos(lat2_rad) y = math.cos(lat1_rad) * math.sin(lat2_rad) - math.sin(lat1_rad) * math.cos(lat2_rad) * math.cos(delta_lon) bearing = math.atan2(x, y) return (math.degrees(bearing) + 360) % 360 # 加载并排序数据(必须确保组内时间有序) df = pl.read_csv("your_flight_data.csv").sort(["plane_no", "timestamp"]) # 添加前一个时间点的GPS坐标 df = df.with_columns( pl.col("gps_lat").shift(1).over("plane_no").alias("prev_lat"), pl.col("gps_lon").shift(1).over("plane_no").alias("prev_lon") ) # 计算方位角(第一行无前置点,结果为None) df = df.with_columns( pl.struct(["prev_lat", "prev_lon", "gps_lat", "gps_lon"]) .map_elements( lambda s: calculate_bearing(*s.values()) if all(v is not None for v in s.values()) else None, return_dtype=pl.Float64 ).alias("bearing") )
2. 优化方案:向量化计算(避免逐行操作)
map_elements是逐行处理,大数据集下性能较差,改用Polars内置的向量化函数实现方位计算:
import polars as pl # 加载并排序数据 df = pl.read_csv("your_flight_data.csv").sort(["plane_no", "timestamp"]) # 添加前序坐标并向量化计算方位角 df = df.with_columns( pl.col("gps_lat").shift(1).over("plane_no").alias("prev_lat"), pl.col("gps_lon").shift(1).over("plane_no").alias("prev_lon") ).with_columns( # 转换为弧度 pl.col("gps_lat").radians().alias("lat_rad"), pl.col("gps_lon").radians().alias("lon_rad"), pl.col("prev_lat").radians().alias("prev_lat_rad"), pl.col("prev_lon").radians().alias("prev_lon_rad") ).with_columns( (pl.col("lon_rad") - pl.col("prev_lon_rad")).alias("delta_lon") ).with_columns( (pl.sin("delta_lon") * pl.cos("lat_rad")).alias("x"), (pl.cos("prev_lat_rad") * pl.sin("lat_rad") - pl.sin("prev_lat_rad") * pl.cos("lat_rad") * pl.cos("delta_lon")).alias("y") ).with_columns( # 计算并格式化方位角,填充第一行的Null ((pl.atan2("x", "y").degrees() + 360) % 360).fill_null(None).alias("bearing") ).drop(["lat_rad", "lon_rad", "prev_lat_rad", "prev_lon_rad", "delta_lon", "x", "y"])
3. 备选方案:分组后apply处理
如果需要对每组进行更复杂的自定义逻辑,可使用group_by().apply():
import polars as pl def process_plane_group(group): # 组内排序(确保时间有序) group = group.sort("timestamp") # 添加前序坐标 group = group.with_columns( pl.col("gps_lat").shift(1).alias("prev_lat"), pl.col("gps_lon").shift(1).alias("prev_lon") ) # 向量化计算方位角 return group.with_columns( pl.when(pl.col("prev_lat").is_not_null()) .then( ((pl.atan2( pl.sin((pl.col("gps_lon") - pl.col("prev_lon")).radians()) * pl.cos(pl.col("gps_lat").radians()), pl.cos(pl.col("prev_lat").radians()) * pl.sin(pl.col("gps_lat").radians()) - pl.sin(pl.col("prev_lat").radians()) * pl.cos(pl.col("gps_lat").radians()) * pl.cos((pl.col("gps_lon") - pl.col("prev_lon")).radians()) ).degrees() + 360) % 360) ) .otherwise(None) .alias("bearing") ) # 加载数据并分组处理 df = pl.read_csv("your_flight_data.csv") df = df.group_by("plane_no").apply(process_plane_group)
关键优化建议
- 强制排序:所有方案都必须先按
plane_no和timestamp排序,否则shift()获取的前序坐标会出错。 - 优先向量化:始终用Polars内置的向量化函数替代逐行操作(如
map_elements),性能可提升数倍甚至数十倍。 - 清理中间列:计算完成后删除不需要的中间列(如弧度列、delta值),减少内存占用。
- 缺失值处理:每组第一行的方位角为Null,可根据业务需求用
fill_null()填充为0、前一个有效值或保留Null。 - 滚动窗口场景适配:如果需要基于最近N个点/时间段计算方位(而非仅前一个点),可使用
group_by().rolling(),示例:# 基于最近1小时的窗口计算(示例,需根据需求调整逻辑) df.group_by("plane_no").rolling(index_column="timestamp", window="1h").agg( # 自定义滚动窗口内的方位计算逻辑 )
内容的提问来源于stack exchange,提问作者jimfawkes
相关产品推荐
相关产品推荐

