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

在线数据摄入阶段Pipeline聚合函数异常问题排查求助

问题排查:MLRun聚合函数在线返回0而非预期6

核心原因1:时间窗口基准与数据时间不匹配

MLRun滑动时间聚合窗口默认以当前系统时间作为基准计算窗口范围。你的测试数据sysdate均为2021-01-01,执行代码时的当前时间远超出该日期+60天的范围,导致所有历史数据被排除在窗口外,最终count结果为0。

解决办法

  • 方案一:修改测试数据的sysdate为接近当前时间的日期,例如:
    current_time = datetime.datetime.now()
    "sysdate": [current_time - datetime.timedelta(days=1) for _ in range(6)]
    
  • 方案二:获取在线特征时指定asof_time,强制以数据中的时间为基准计算窗口:
    resp = svc.get([{"key0": 1, "key1":0} ], asof_time=datetime.datetime(2021,1,1,2))
    

核心原因2:Feature Set未配置时间戳列

代码中未明确指定Feature Set的时间基准列(timestamp_key),MLRun无法正确识别sysdate作为聚合的时间依据,导致窗口计算逻辑错误。

解决办法

创建Feature Set时显式配置时间列:

feature_set = fstore.get_or_create_feature_set(
    "sample", 
    project=project.name,
    entities=[fstore.Entity("key0"), fstore.Entity("key1")],
    timestamp_key="sysdate",
    targets=[NoSqlTarget(), ParquetTarget()]
)

其他潜在问题及修正

  1. 未定义变量报错:代码中featureGetOrCreate、input_df未定义,实际运行会触发异常,需补充:
    import pandas as pd
    input_df = pd.DataFrame(data)
    # 替换自定义函数为MLRun原生API
    feature_set = fstore.get_or_create_feature_set(...)
    
  2. 在线聚合计算触发:确保fstore.ingest正确执行了聚合逻辑,可通过检查output_df确认聚合结果是否生成,再同步到在线存储。

验证修改后的完整代码示例

import datetime
import pandas as pd
import mlrun
import mlrun.feature_store as fstore
from mlrun.datastore.targets import ParquetTarget, NoSqlTarget

# 准备数据,sysdate设为近期日期
current_time = datetime.datetime.now()
data = {
    "key0": [1,1,1,1,1,1], 
    "key1": [0,0,0,0,0,0],
    "fn1": [1,1,2,3,1,0],
    "sysdate": [current_time - datetime.timedelta(days=1) for _ in range(6)]
}
input_df = pd.DataFrame(data)

# 创建项目与Feature Set,明确时间列和存储目标
project = mlrun.get_or_create_project("jist-agg", context='./', user_project=False)
feature_set = fstore.get_or_create_feature_set(
    "sample", 
    project=project.name,
    entities=[fstore.Entity("key0"), fstore.Entity("key1")],
    timestamp_key="sysdate",
    targets=[NoSqlTarget(), ParquetTarget()]
)

# 添加60天窗口count聚合
feature_set.add_aggregation(name='fn1', column='fn1', operations=['count'], windows=['60d'], step_name="agg1")

# 摄入数据到在线/离线存储
output_df = fstore.ingest(feature_set, input_df, overwrite=True, infer_options=fstore.InferOptions.default())

# 读取在线特征并验证
svc = fstore.get_online_feature_service(fstore.FeatureVector("my-vec", ["sample.*"], with_indexes=True))
resp = svc.get([{"key0": 1, "key1":0} ])
assert resp[0]['fn1_count_60d'] == 6.0, 'Mistake in solution'

内容的提问来源于stack exchange,提问作者JIST

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 14:40:24