工业 AI 在某个时点判断设备状态时,不同传感器可能没有恰好同一秒的记录。本例按 device_id 找该时点之前已经到达的最近读数,再检查测量本身是否过旧。这样能避免将后到的数据错误塞进历史训练样本。
设备、时间、温度和到达延迟均为人工构造;已在本地运行融合代码,没有连接工业设备、训练模型或执行控制指令。本题是工业数据融合智能体外围工具的确定性处理步骤,模型不能靠猜测补齐传感器记录。

区分测量时间、到达时间与判断时点
| 字段 | 含义 | 为什么需要 |
|---|---|---|
| observation_at | 要构建特征或判断状态的时点 | 决定当时能使用哪些信息 |
| reading_at | 传感器实际测量时间 | 判断观测是否足够新鲜 |
| available_at | 此读数在处理系统中可用的时间 | 防止已测量但尚未到达的数据被提前使用 |
| device_id / reading_id | 设备与读数身份 | 避免串设备,允许回查具体原记录 |
a05 虽在 10:05 测量,但直到 10:11 才可用;o03 的时点是 10:10,不能使用这条尚未到达的读数。只按 reading_at 做历史匹配会忽略这项延迟。本例两种时间都已明确为 UTC,时钟偏差与来源时间协议应在接入前另行校验。
匹配已到达记录,再做新鲜度检查
第一道规则:按 available_at 升序,以 backward 选择不晚于 observation_at 的同设备读数,并限制到达时间差不超过 6 分钟。第二道规则:该读数的 reading_at 距判断时点也不能超过 6 分钟。前者防止提前使用,后者避免刚到达的旧测量被当作新状态。
pandas merge_asof 要求两侧按连接时间键全局升序;不能先按设备分组排序,再让时间在设备之间倒退。by=”device_id” 负责匹配组内关系,代码另外拒绝重复的设备/到达时间组合,以免相同时间有多个候选时不知应选谁。
示例实际运行环境为 Windows、Python 3.11.15、pandas 2.3.3。
在空的练习目录里创建独立环境。以下为 Windows PowerShell 命令;macOS/Linux 的激活命令用 source .venv/bin/activate。输出文件只用于本次演示,重复执行会更新这些演示文件。
python -m venv .venv
.\.venv\Scripts\Activate.ps1
python -m pip install pandas==2.3.3
保存代码为 industrial_time_fusion.py,执行 python industrial_time_fusion.py,输出 fusion_output/fusion_audit.csv。
from pathlib import Path
import pandas as pd
# Artificial observations and readings; all timestamps already use UTC.
observations = pd.DataFrame({
'observation_id': ['o01', 'o02', 'o03', 'o04', 'o05', 'o06', 'o07'],
'device_id': ['A', 'B', 'A', 'B', 'A', 'C', 'A'],
'observation_at': pd.to_datetime([
'2026-10-01T10:03:00Z', '2026-10-01T10:04:00Z',
'2026-10-01T10:10:00Z', '2026-10-01T10:12:00Z',
'2026-10-01T10:30:00Z', '2026-10-01T10:10:00Z',
'2026-10-01T10:11:30Z',
], utc=True),
})
readings = pd.DataFrame({
'reading_id': ['a00', 'b02', 'a05', 'b09', 'a20'],
'device_id': ['A', 'B', 'A', 'B', 'A'],
'reading_at': pd.to_datetime([
'2026-10-01T10:00:00Z', '2026-10-01T10:02:00Z',
'2026-10-01T10:05:00Z', '2026-10-01T10:09:00Z',
'2026-10-01T10:20:00Z',
], utc=True),
'available_at': pd.to_datetime([
'2026-10-01T10:00:30Z', '2026-10-01T10:02:30Z',
'2026-10-01T10:11:00Z', '2026-10-01T10:09:30Z',
'2026-10-01T10:20:30Z',
], utc=True),
'temperature_c': [20, 40, 21, 41, 24],
})
assert observations['observation_id'].is_unique
assert readings['reading_id'].is_unique
assert not readings.duplicated(['device_id', 'available_at']).any()
assert readings[['reading_at', 'available_at']].notna().all().all()
assert readings['reading_at'].le(readings['available_at']).all()
observations['input_order'] = range(len(observations))
max_age = pd.Timedelta('6min')
joined = pd.merge_asof(
observations.sort_values('observation_at'),
readings.sort_values('available_at'),
left_on='observation_at', right_on='available_at',
by='device_id', direction='backward', tolerance=max_age,
).sort_values('input_order')
has_reading = joined['reading_id'].notna()
joined['availability_age_seconds'] = (
joined['observation_at'] - joined['available_at']
).dt.total_seconds()
joined['measurement_age_seconds'] = (
joined['observation_at'] - joined['reading_at']
).dt.total_seconds()
fresh = has_reading & joined['measurement_age_seconds'].le(max_age.total_seconds())
joined['status'] = 'no_recent_available_reading'
joined.loc[~joined['device_id'].isin(readings['device_id']), 'status'] = 'no_device_reading'
joined.loc[has_reading & ~fresh, 'status'] = 'stale_measurement'
joined.loc[fresh, 'status'] = 'matched'
# Keep the raw candidate in the audit table; only accepted values enter a feature.
joined['feature_temperature_c'] = joined['temperature_c'].where(fresh)
assert joined.loc[has_reading, 'available_at'].le(joined.loc[has_reading, 'observation_at']).all()
assert joined.loc[has_reading, 'reading_at'].le(joined.loc[has_reading, 'observation_at']).all()
expected = ['matched', 'matched', 'no_recent_available_reading', 'matched',
'no_recent_available_reading', 'no_device_reading', 'stale_measurement']
assert joined['status'].tolist() == expected
assert joined.loc[fresh, 'reading_id'].tolist() == ['a00', 'b02', 'b09']
out = Path('fusion_output')
out.mkdir(exist_ok=True)
joined.drop(columns='input_order').to_csv(out / 'fusion_audit.csv', index=False, encoding='utf-8-sig')
print(joined[['observation_id', 'device_id', 'reading_id', 'feature_temperature_c', 'status']].to_string(index=False))
print('observations:', len(joined), 'accepted:', int(fresh.sum()))
print('future_available_readings_used:', 0)
按观测编号读结果
observation_id device_id reading_id feature_temperature_c status
o01 A a00 20.0 matched
o02 B b02 40.0 matched
o03 A NaN NaN no_recent_available_reading
o04 B b09 41.0 matched
o05 A NaN NaN no_recent_available_reading
o06 C NaN NaN no_device_reading
o07 A a05 NaN stale_measurement
observations: 7 accepted: 3
future_available_readings_used: 0
七条观测保留七行,其中 o01、o02、o04 接受的温度分别是 20、40、41℃;其余四条的特征值保持缺失,并有具体状态。o03 尚不能用 a05,且更早的 a00 已超到达时间范围;o05 最近可用读数也超时;o06 的设备 C 根本没有读数。
o07 在 10:11:30 能看到刚到达的 a05,但该测量来自 10:05,已经过了 6 分 30 秒,所以状态为 stale_measurement。代码保留候选读数供回查,但只有 status 为 matched 的值进入 feature_temperature_c;不能绕过状态直接把 temperature_c 全列送进模型。
打开 CSV,核对 matched 行的 device_id、reading_id、reading_at 与 available_at 是否来自原读数,并查看 availability_age_seconds 和 measurement_age_seconds。前者分别为 150、90、150 秒,后者分别为 180、120、180 秒;所有接受的来源时间都不晚于判断时点。
把这一步接到工业 AI 流程
受限融合工具输入应指定观测批次、设备列和已批准的新鲜度规则;输出融合表以及未匹配清单。智能体可以解释哪些记录需要补采或回查,但不能把未匹配温度默认填 0、凭相邻设备读数代替,或自动向设备下发控制命令。
6 分钟只是本例教学阈值,真实设备的新鲜度必须由采样协议、数据延迟和模型任务确定,不能机械照抄。若测量流会乱序到达,本例定义为“最近到达的候选再检查新鲜度”;需要“所有已到达读数中测量时间最新”的策略时,应另写并验证该选择规则。
真正用于历史训练时,还需要保存当时可用的记录版本。今天才修正的历史读数即使 reading_at 很早,也不能默认表示当年的模型已经知道;available_at 应反映真实可用时间,而不是从测量时间复制。完成融合后再按任务划分训练与验证,本文不提供预测准确率结论。
官方资料
以下资料实际读取于 2026 年 10 月 1 日;接口或界面范围按正文注明的版本核对。
Ai菜鸟网。发布者:AI小管家,转载请注明出处:https://www.alyyhw.com/32727.html