如何用Polars LazyFrames实现时间戳同步(替代Numpy实现)
用Polars LazyFrames实现时间戳同步(替代Numpy方案)
需求说明
需要基于Polars LazyFrames实现两个时间戳数组的同步功能,替代原有Numpy实现以处理大规模数据。核心逻辑是在指定容差内,为每个时间戳匹配另一个数组中最近的时间戳,返回符合要求的时间戳对。
现有数据与Numpy实现
数据定义(Polars LazyFrames)
import polars as pl import numpy as np timestamps = pl.LazyFrame( np.array( [ np.datetime64("1970-01-01T00:00:00.500000000"), np.datetime64("1970-01-01T00:00:01.500000000"), np.datetime64("1970-01-01T00:00:02.600000000"), np.datetime64("1970-01-01T00:00:03.400000000"), np.datetime64("1970-01-01T00:00:04.500000000"), np.datetime64("1970-01-01T00:00:05.300000000"), np.datetime64("1970-01-01T00:00:06.200000000"), np.datetime64("1970-01-01T00:00:07.400000000"), np.datetime64("1970-01-01T00:00:08.500000000"), ] ), schema={"values": pl.Datetime} ) other_timestamps = pl.LazyFrame( np.array( [ np.datetime64("1970-01-01T00:00:01.500000000"), np.datetime64("1970-01-01T00:00:02.000000000"), np.datetime64("1970-01-01T00:00:02.500000000"), np.datetime64("1970-01-01T00:00:04.500000000"), np.datetime64("1970-01-01T00:00:06.000000000"), np.datetime64("1970-01-01T00:00:06.500000000"), ] ), schema={"values": pl.Datetime} )
Numpy实现代码
import numpy as np import numpy.typing as npt def parse_timedelta(tolerance_str: str) -> np.timedelta64: if tolerance_str.endswith("ms"): return np.timedelta64(int(tolerance_str[:-2]), "ms") elif tolerance_str.endswith("s"): return np.timedelta64(int(tolerance_str[:-1]), "s") else: raise ValueError("Unsupported tolerance format") def _np_sync_to( timestamps: npt.ArrayLike[np.datetime64], other: npt.ArrayLike[np.datetime64], tolerance: str, ): outer_diffs = np.abs(np.subtract.outer(other, timestamps)) closest_timestamps_indices = outer_diffs.argmin(0) closest_timestamps = other[closest_timestamps_indices] diffs = np.abs(closest_timestamps - timestamps) tolerance = parse_timedelta(tolerance) within_tolerance = diffs <= tolerance ts1_synced = timestamps[within_tolerance] ts2_synced = closest_timestamps[within_tolerance] return ts1_synced, ts2_synced np_ts1_synced, np_ts2_synced = _np_sync_to( timestamps=np.squeeze(timestamps.collect().to_numpy()), other=np.squeeze(other_timestamps.collect().to_numpy()), tolerance="500ms", )
预期结果
np_ts1_synced = np.array([ np.datetime64('1970-01-01T00:00:01.500000000'), np.datetime64('1970-01-01T00:00:02.600000000'), np.datetime64('1970-01-01T00:00:04.500000000'), np.datetime64('1970-01-01T00:00:06.200000000') ]) np_ts2_synced = np.array([ np.datetime64('1970-01-01T00:00:01.500000000'), np.datetime64('1970-01-01T00:00:02.500000000'), np.datetime64('1970-01-01T00:00:04.500000000'), np.datetime64('1970-01-01T00:00:06.000000000') ])
问题
尝试了两种Polars实现方式均未得到预期结果:
- 模拟
numpy.subtract.outer的方法存在维度计算错误; - 直接使用
join_asof仅能单向匹配,无法获取最近的时间戳。
正确的Polars LazyFrames实现方案
通过双向join_asof+ 差值比较的方式实现,既保证LazyFrame的高效性,又能准确匹配最近时间戳:
import polars as pl # 1. 定义容差 tolerance = pl.parse_duration("500ms") # 2. 重命名other_timestamps的列,避免合并时冲突 other_renamed = other_timestamps.rename({"values": "other_values"}) # 3. 分别执行向前、向后的asof join # 向前匹配:取小于等于当前时间的最近other时间戳 forward_join = timestamps.join_asof( other_renamed, left_on="values", right_on="other_values", direction="forward" ) # 向后匹配:取大于等于当前时间的最近other时间戳 backward_join = timestamps.join_asof( other_renamed, left_on="values", right_on="other_values", direction="backward" ) # 4. 合并两个方向的结果,计算差值并选择最近的时间戳 synced_result = forward_join.join( backward_join, on="values", suffix="_backward" ).with_columns( # 计算两个方向的时间差绝对值 diff_forward=pl.abs(pl.col("values") - pl.col("other_values")), diff_backward=pl.abs(pl.col("values") - pl.col("other_values_backward")), # 选择差值更小的other时间戳 closest_other=pl.when(pl.col("diff_forward") <= pl.col("diff_backward")) .then(pl.col("other_values")) .otherwise(pl.col("other_values_backward")), # 保留最小差值用于过滤 min_diff=pl.min([pl.col("diff_forward"), pl.col("diff_backward")]) ).filter( # 过滤出容差范围内的记录 pl.col("min_diff") <= tolerance ).select( # 提取最终的时间戳对 ts1_synced=pl.col("values"), ts2_synced=pl.col("closest_other") ) # 5. 执行LazyFrame计划并获取结果 result = synced_result.collect() ts1_synced = result["ts1_synced"].to_numpy() ts2_synced = result["ts2_synced"].to_numpy() # 验证结果 print("ts1_synced:\n", ts1_synced) print("\nts2_synced:\n", ts2_synced)
方案说明
- 双向
join_asof:分别获取每个时间戳前后最近的候选值,避免单向匹配的局限性; - 差值比较:通过计算两个候选值与当前时间戳的差值,选择最近的那个;
- 容差过滤:仅保留差值在指定范围内的时间戳对;
- LazyFrame优化:所有操作均基于LazyFrame执行,Polars会自动优化查询计划,适合处理大规模数据集。
内容的提问来源于stack exchange,提问作者jspieler
相关产品推荐
相关产品推荐

