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

使用ThreadPoolExecutor并行处理Keras模型时结果异常问题

Keras模型ThreadPoolExecutor并行预测结果错误问题

问题现象

使用concurrent.futures.ThreadPoolExecutor并行执行Keras模型预测(针对不同权重集计算误差)时,输出结果与串行处理不一致,且每次并行运行结果存在随机性。即使设置了NumPy随机种子,问题仍未解决。

可复现代码

import tensorflow.keras
import numpy
import concurrent.futures

numpy.random.seed(1)

def create_rand_weights(model, num_models):
    random_model_weights = []
    for model_idx in range(num_models):
        random_weights = []
        for layer_idx in range(len(model.weights)):
            layer_shape = model.weights[layer_idx].shape
            if len(layer_shape) > 1:
                layer_weights = numpy.random.rand(layer_shape[0], layer_shape[1])
            else:
                layer_weights = numpy.random.rand(layer_shape[0])
            random_weights.append(layer_weights)
        random_weights = numpy.array(random_weights, dtype=object)
        random_model_weights.append(random_weights)
    
    random_model_weights = numpy.array(random_model_weights)
    return random_model_weights

def model_error(model_weights):
    global data_inputs, data_outputs, model
    model.set_weights(model_weights)
    predictions = model.predict(data_inputs)
    mae = tensorflow.keras.losses.MeanAbsoluteError()
    abs_error = mae(data_outputs, predictions).numpy() + 0.00000001
    return abs_error

input_layer  = tensorflow.keras.layers.Input(3)
dense_layer1 = tensorflow.keras.layers.Dense(5, activation="relu")(input_layer)
output_layer = tensorflow.keras.layers.Dense(1, activation="linear")(dense_layer1)
model = tensorflow.keras.Model(inputs=input_layer, outputs=output_layer)

data_inputs = numpy.array([[0.02, 0.1, 0.15],
                           [0.7, 0.6, 0.8],
                           [1.5, 1.2, 1.7],
                           [3.2, 2.9, 3.1]])    
data_outputs = numpy.array([[0.1],
                            [0.6],
                            [1.3],
                            [2.5]])

num_models = 10
random_model_weights = create_rand_weights(model, num_models)

ExecutorClass = concurrent.futures.ThreadPoolExecutor
thread_output = []
with ExecutorClass(max_workers=2) as executor:
    output = executor.map(model_error, random_model_weights)
for out in output:
    thread_output.append(out)
thread_output=numpy.array(thread_output)
print("Wrong Outputs using Threads")
print(thread_output)

print("\n\n")

correct_output = []
for idx in range(num_models):
    error = model_error(random_model_weights[idx])
    correct_output.append(error)
correct_output=numpy.array(correct_output)
print("Correct Outputs without Threads")
print(correct_output)

输出对比

  • 串行处理正确输出:
[6.78012372 3.42922212 4.96738673 6.64474774 6.83102609 4.41165734 3.34482099 7.6132908  7.97145654 6.98378612]
  • 并行处理错误输出(每次运行结果可能不同):
[3.42922212 3.42922212 6.90911246 6.64474774 4.41165734 3.34482099 7.6132908  7.97145654 6.98378612 6.98378612]

问题原因

Keras/TensorFlow的模型对象不具备线程安全性。多个线程同时调用model.set_weights()和model.predict()时,会出现竞态条件:线程A刚设置好权重,线程B就覆盖了该权重,导致线程A执行预测时使用的是线程B的权重,最终输出错误结果。

解决方案

方案1:为每个线程创建独立模型实例

在每个任务内部重新构建模型并加载对应权重,避免共享同一个模型对象。修改model_error函数如下:

def model_error(model_weights):
    global data_inputs, data_outputs
    # 每个线程独立创建模型
    input_layer  = tensorflow.keras.layers.Input(3)
    dense_layer1 = tensorflow.keras.layers.Dense(5, activation="relu")(input_layer)
    output_layer = tensorflow.keras.layers.Dense(1, activation="linear")(dense_layer1)
    model = tensorflow.keras.Model(inputs=input_layer, outputs=output_layer)
    
    model.set_weights(model_weights)
    predictions = model.predict(data_inputs)
    mae = tensorflow.keras.losses.MeanAbsoluteError()
    abs_error = mae(data_outputs, predictions).numpy() + 0.00000001
    return abs_error

方案2:使用线程锁同步模型访问

通过全局锁确保同一时间只有一个线程操作模型,避免竞态条件。但这种方法会让并行处理退化为串行,失去并行提速的意义,仅适合必须共享模型的场景:

import threading

# 定义全局锁
model_lock = threading.Lock()

def model_error(model_weights):
    global data_inputs, data_outputs, model
    with model_lock:
        model.set_weights(model_weights)
        predictions = model.predict(data_inputs)
    mae = tensorflow.keras.losses.MeanAbsoluteError()
    abs_error = mae(data_outputs, predictions).numpy() + 0.00000001
    return abs_error

方案3:改用ProcessPoolExecutor

多进程环境中,每个进程拥有独立的模型副本,不会出现权重共享冲突。只需修改Executor类:

ExecutorClass = concurrent.futures.ProcessPoolExecutor

注意:多进程会带来更高的内存开销,因为每个进程都会加载模型和数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 16:17:56