在线数据摄入阶段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()] )
其他潜在问题及修正
- 未定义变量报错:代码中
featureGetOrCreate、input_df未定义,实际运行会触发异常,需补充:import pandas as pd input_df = pd.DataFrame(data) # 替换自定义函数为MLRun原生API feature_set = fstore.get_or_create_feature_set(...) - 在线聚合计算触发:确保
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
相关产品推荐
相关产品推荐

