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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 05:55:54