joblib Parallel设置n_jobs>1时运行结果与串行不一致问题
Joblib Parallel并行执行嵌套交叉验证结果与串行不一致问题
我特意强调“不正确”,是因为问题并非Parallel()运行速度慢于或等同于串行代码,而是输出结果存在本质差异。
我在Google Colab环境下开发,尝试对for循环做并行化(具体为嵌套交叉验证循环的外层循环)。
代码逻辑为:外层for循环沿数据表滑动切分窗口,内层for循环遍历不同参数训练模型、筛选最优参数,随后使用全部“前置”数据训练最终模型,在未触碰的测试折上完成预测。
简化示例代码
from joblib import Parallel, delayed def NestedCV(outer): traindataforthisiteration = data[0:(1000+outer*500),:] testdataforthisiteration = data[(1000+outer*500):(1000+(outer+1)*500),:] bestresult = 0 bestparam = 0 paramoptions = [1,2,3] for i in paramoptions: # 遍历所有参数选项,记录结果,遇到更优结果时更新bestparam和bestresult innertrain = traindataforthisiteration[0:int((1000+outer*500)/2),:] innertest = traindataforthisiteration[int((1000+outer*500)/2):(1000+outer*500),:] innermodel = model(i).fit(innertrain) # 用训练好的模型在内部测试集上预测,得到性能指标 performance = accuracy(model.predict(innertest)) if performance > bestresult: bestresult = performance bestparam = i finalmodel = model(bestparam).fit(traindataforthisiteration) results = finalmodel.predict(testdataforthisiteration) return results, bestparam
对应的并行调用代码如下:
inputs = range(3) results = Parallel(n_jobs=3)(delayed(NestedCV)(i) for i in inputs)
对照测试结果
上述并行代码的输出,与串行运行得到的(推测为正确的)结果存在差异。已完成以下三组对照测试,三组测试输出完全一致,但均与n_jobs=3并行运行的结果不同:
- 复制
NestedCV函数内的代码,不做函数封装,设置outer=0直接运行; - 按顺序单独调用
NestedCV(0)、NestedCV(1)、NestedCV(2); - 设置
n_jobs=1调用Parallel,即强制使用串行执行模式。
补充信息
- 我不常使用Google Colab,也很少编写并行代码,示例中的
Parallel(n_jobs=3)(delayed(...))写法参考自公开博客,因此错误可能属于非常基础的用法问题,但暂时未定位到根因。我曾怀疑是循环迭代间存在依赖导致并行失效,但如果是该原因,随机顺序调用NestedCV(n)也应该得到不一致的结果,实际测试并未出现该情况。 - 我的函数存在多个返回值,不确定是否会影响结果、是否存在输出捕获错误的可能,多返回值的设计已经在上述函数代码中体现。
- 实际使用的代码为基于Stable-Baselines3编写的强化学习智能体嵌套交叉验证代码,简化版如下:
def NestedCV(outer): lrparams = [0.0003, 0.0001] batchsizes = [256] gammaoptions = [0.9] starttimestep = 250000 increment = 100000 # 注意inctimes的遍历范围,当rep==0时训练步长为starttimestep inctimes = 1 outertrainltc = ltctrain.iloc[0:(outercvstart+outer*outerfoldsize),:] outertrainpriceltc = ltcpricestrain[0:(outercvstart+outer*outerfoldsize)] outertest = train.iloc[((outercvstart+outer*outerfoldsize)+1):(outercvstart+(outer+1)*outerfoldsize),:] outertestprice = pricestrain[((outercvstart+outer*outerfoldsize)+1):(outercvstart+(outer+1)*outerfoldsize)] bestlr = 0 bestf = 0 bestr = 0 bestg = 0 bestrew = -np.inf for lr in lrparams: for r in batchsizes: for g in gammaoptions: innerres1 = [] innerres2 = [] for inner in range(2): lastind = (len(outertrain.index)-3*innerfoldsize+inner*innerfoldsize) # 训练模型第一阶段(规则学习阶段) env = StockPredEnv(df1=outertrainltc.iloc[0:lastind,:], priceseries1=list(outertrainpriceltc[0:lastind]), episodelen= 2225, totaltimesteps=starttimestep) env = ActionMasker(env, mask_fn) model = MaskablePPO(MaskableActorCriticPolicy, env, verbose=0, learning_rate=lr, gamma=g, batch_size=r, seed=10) for rep in range(inctimes): if rep==0: model.learn(starttimestep) else: model.learn(increment) # 在测试数据上验证效果 # 切换为测试数据对应的环境 env = StockPredEnv(df1=outertrainltc.iloc[lastind:(lastind+innerfoldsize),:], priceseries1=list(outertrainpriceltc[lastind:(lastind+innerfoldsize)]), episodelen= len(list(outertrainprice[lastind:(lastind+innerfoldsize)]))-1, totaltimesteps=starttimestep+rep*increment) env = ActionMasker(env, mask_fn) obs = env.reset().astype('float') done = False score = 0 while not done: action_masks = get_action_masks(env) action, _ = model.predict(obs, action_masks=action_masks, deterministic=True) obs, reward, done, info = env.step(action[0]) score += reward if inner==0: innerres1.append(score) else: innerres2.append(score) # 计算每个时间步的平均得分 meanscores = [] for m in range(len(innerres1)): meanscores.append(np.mean([innerres1[m], innerres2[m]])) # 计算最优结果和对应的训练步长 topmean = -np.inf toptimesteps = 0 for m in range(len(meanscores)): if meanscores[m] > topmean: topmean = meanscores[m] toptimesteps = starttimestep + m*increment # 和历史最优结果对比,记录更优的超参数 if topmean > bestrew: bestrew = topmean bestlr = lr bestf = toptimesteps bestr = r bestg = g # 记录最优参数 bestparams = (bestlr, bestf, bestr, bestg) cvbestrew = bestrew # 用全部内部训练数据训练最终模型,在外层测试折上验证 # 训练模型第一阶段(规则学习阶段) env = StockPredEnv(df1=outertrainltc, priceseries1=list(outertrainpriceltc), episodelen= 2225, totaltimesteps=bestf) env = ActionMasker(env, mask_fn) model = MaskablePPO(MaskableActorCriticPolicy, env, verbose=0, learning_rate=bestlr, gamma=bestg, batch_size=bestr, seed=10) model.learn(bestf) # 在测试数据上验证 # 切换为包含测试数据的环境 env = StockPredEnv(df1=outertrainltc, priceseries1=list(outertrainpriceltc), episodelen= len(list(outertestprice))-1, totaltimesteps=len(list(outertestprice))) env = ActionMasker(env, mask_fn) lastepisodeactions = [] obs = env.reset().astype('float') done = False score = 0 while not done: action_masks = get_action_masks(env) action, _ = model.predict(obs, action_masks=action_masks, deterministic=True) obs, reward, done, info = env.step(action[0]) score += reward lastepisodeactions.append(action) outerfoldrew = score tradeactions = lastepisodeactions return bestparams, outerfoldrew, tradeactions
- 我曾尝试用完全相同的输入并行运行3次函数实例,3次并行输出完全一致,但仍然和串行运行结果不同。这说明问题并非竞态条件导致——否则3次并行运行的结果应该互不相同,希望了解该现象对应的可能原因与排查方向。
内容的提问来源于stack exchange,提问作者Vladimir Belik
相关产品推荐
相关产品推荐

