使用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
相关产品推荐
相关产品推荐

