如何在超过1.2亿行的数据中,使用带有最大时间间隔(GAP)的前向填充(ffill)?
问题:
我在一个超过1.2亿行的数据集上遇到了性能瓶颈。
我需要对传感器数据中的NaN值进行向前填充,但只有当距离上一次有效读取的时间间隔小于5分钟时才进行填充。
如果间隔更大,则必须保持为NaN。
A standard groupby().apply() takes over 40 minutes, which is too slow for our pipeline.
Oops, I mistakenly left English sentence here. We must translate all content. The line "A standard groupby().apply() takes over 40 minutes, which is too slow for our pipeline." should be translated. Let's correct.
The translation for that line:
标准的 groupby().apply() 需要超过40分钟,这对我们的流水线来说太慢。
Let's re-output the full translated content cleanly:
问题:
我在一个超过1.2亿行的数据集上遇到了性能瓶颈。
我需要对传感器数据中的NaN值进行向前填充,但只有当距离上一次有效读取的时间间隔小于5分钟时才进行填充。
如果间隔更大,则必须保持为NaN。
标准的 groupby().apply() 需要超过40分钟,这对我们的流水线来说太慢。
我的代码:
import pandas as pd
import numpy as np
df = pd.DataFrame({
'sensor_id': [1, 1, 1, 1, 2, 2],
'timestamp': pd.to_datetime(['10:00', '10:03', '10:10', '10:11', '10:00', '10:01']),
'temp': [22.1, np.nan, np.nan, 23.0, 19.5, np.nan]
})
# This works on small data but is painfully slow on 100M rows
def limit_ffill(group):
return group['temp'].ffill(limit=1)
# df['temp'] = df.groupby('sensor_id').apply(limit_ffill)
目标:
我需要一种向量化的方法(或者也许使用numpy或numba)来实现。
本质上:
- 以
sensor_id分组。 - 向前填充
temp。 - 如果
current_timestamp - last_valid_timestamp > 5 minutes,则重置为NaN。
有没有办法在没有Python循环或笨重的 .groupby().apply() 的情况下做到?
有没有一种跨组在 timedelta 限制下的向量化方式来 ffill?
我看过 pd.merge_asof,但我还不太清楚如何让它用于一个简单的向前填充。
解决方案
我需要一种向量化的方法(或许使用
numpy或numba)来实现。
我个人建议使用Numba——通常这是获得这种特定行为的最简单方法。
我会给出如下思路的算法:
- 按时间戳排序数据。
- 遍历行,记录每个传感器的最后一次观测。
- 如果观测值是nan,在最近时间足够近时进行前向填充。
- 如果不是nan,则将其记录为一个新的观测。
这使你能够在一次遍历数据时完成迭代步骤。
代码
import pandas as pd
import numpy as np
import numba as nb
import datetime
df = pd.DataFrame({
'sensor_id': [1, 1, 1, 1, 2, 2],
'timestamp': pd.to_datetime(['10:00', '10:03', '10:10', '10:11', '10:00', '10:01']),
'temp': [22.1, np.nan, np.nan, 23.0, 19.5, np.nan]
})
def ffill_in_group_ignoring_gaps(df, max_gap_s):
# Note: algorithm below assumes dataset is sorted by timestamp
# Could skip this if you know this property already holds
df = df.sort_values('timestamp')
group_col, timestamp_col, x_col = df['sensor_id'], df['timestamp'], df['temp']
# sensor_id might be sparsely populated, e.g. have sensor values of 1, 100, 101.
# Make this dense with factorize()
sensor_id_codes, sensor_ids_uniques = pd.factorize(group_col)
num_sensors = len(sensor_ids_uniques)
# Convert timestamp to seconds since epoch
# You can use .astype('int64') directly if you know the unit of your timestamp
timestamp_epoch = (timestamp_col - pd.Timestamp("1970-01-01")) // pd.Timedelta("1s")
# Convert arguments to known types
timestamp_epoch = timestamp_epoch.values.astype('int64')
last_sensor_reading_ts = np.zeros(num_sensors, dtype='int64')
last_sensor_reading_x = np.full(num_sensors, np.nan, dtype='float64')
x_col = x_col.values.astype('float64')
new_temp = ffill_in_group_ignoring_gaps_inner(
sensor_id_codes,
x_col,
last_sensor_reading_ts,
last_sensor_reading_x,
timestamp_epoch,
max_gap_s
)
df['temp'] = new_temp
# Optionally you can use sort_index to restore the original order
return df
@nb.njit()
def ffill_in_group_ignoring_gaps_inner(
sensor_id_codes,
x_col,
last_sensor_reading_ts,
last_sensor_reading_x,
timestamp_epoch,
max_gap_s
):
N = len(sensor_id_codes)
ret = np.zeros(N, dtype='float64')
for i in range(N):
# Read inputs for this row
sensor = sensor_id_codes[i]
now = timestamp_epoch[i]
x_val = x_col[i]
last_reading_ts = last_sensor_reading_ts[sensor]
ffill_val = last_sensor_reading_x[sensor]
if not (now - last_reading_ts < max_gap_s):
# If longer than max_gap_s, don't use last reading
ffill_val = np.nan
if np.isnan(x_val):
# If x=nan, use the last reading.
x_val = ffill_val
else:
# If x!=nan, then write to the last reading array, updating
# x and timestamp
last_sensor_reading_x[sensor] = x_val
last_sensor_reading_ts[sensor] = now
ret[i] = x_val
return ret
print(ffill_in_group_ignoring_gaps(df.copy(), max_gap_s=5*60))
我发现这种方法在以下假设下大约快20倍:
- 数据有5000万行,每个传感器有100次读数。
- 你不能保证输入数据按时间戳排序,因此需要
sort_values()。 (如果按sensor_id排序再timestamp,这也能工作。) 这一步占算法时间的40%。 - 你不关心输出行的顺序。 (代码中有关于如何用
sort_index()来恢复顺序的注释。) - 这并不假设
sensor_id是整数,或它不是密集分布的。比如,你可能有一个sensor_id远大于该数字。pd.factorize()可以处理此类情况,但会耗去大约15% 的时间。 - 我在Pandas 2.2.3与 Numba 0.61.2下进行了基准测试。
以下是我使用的基准。
基准测试
def limit_ffill(group):
return group['temp'].ffill(limit=1)
def ffill_slow(df):
df = df.copy()
df['temp'] = df.groupby('sensor_id').apply(limit_ffill).values
return df
def gen_data(num_readings):
# Source - https://stackoverflow.com/a/57520939
# Posted by Fabian Bosler
# Retrieved 2026-03-20, License - CC BY-SA 4.0
readings_per_sensor = 100
num_sensors = num_readings // readings_per_sensor
N = num_sensors * readings_per_sensor
sensors = np.tile(np.arange(num_sensors), readings_per_sensor)
now = datetime.datetime.now().timestamp()
timestamp = pd.to_datetime(np.linspace(now, now + 4 * 60 * readings_per_sensor, N, endpoint=False), unit='s')
temp = np.random.rand(N)
mask = np.random.rand(N) < 0.5
temp[mask] = np.nan
dft = pd.DataFrame({
'sensor_id': sensors,
'timestamp': timestamp,
'temp': temp,
})
return dft
dft = gen_data(50_000_000)
%timeit -n 1 -r 1 ffill_in_group_ignoring_gaps(dft, max_gap_s=5*60)
%timeit -n 1 -r 1 ffill_slow(dft)
结果:
4.44 s ± 0 ns per loop (mean ± std. dev. of 1 run, 1 loop each)
1min 28s ± 0 ns per loop (mean ± std. dev. of 1 run, 1 loop each)