TensorFlow分布式训练(NCCL+MPI)出现Segmentation Fault求助
问题描述
我正在配备8台RTX2080 GPU节点的集群上,针对时序数据训练简单前馈神经网络。借助MPI完成集群节点组网,采用NCCL作为TensorFlow分布式通信策略,但在模型拟合阶段出现Segmentation Fault(段错误)。我是TensorFlow新手,难以定位问题根源,恳请专业指导。
错误日志
WARNING:tensorflow:`period` argument is deprecated. Please use `save_freq` to specify the frequency in number of batches seen. 2023-01-18 16:23:02.085096: I tensorflow/compiler/mlir/mlir_graph_optimization_pass.cc:185] None of the MLIR Optimization Passes are enabled (registered 2) [g03:35503:0:35774] Caught signal 11 (Segmentation fault: address not mapped to object at address (nil)) ==== backtrace (tid: 35774) ==== 0 0x000000000002137e ucs_debug_print_backtrace() /umbc/ebuild-soft/cascade-lake/build/UCX/1.10.0/GCCcore-10.3.0/ucx-1.10.0/src/ucs/debug/debug.c:656 1 0x000000000382045b tensorflow::NcclCommunicator::Enqueue() collective_communicator.cc:0 2 0x0000000005c9f88a tensorflow::NcclReducer::Run() ???:0 3 0x00000000009086dc tensorflow::BaseCollectiveExecutor::ExecuteAsync(tensorflow::OpKernelContext*, tensorflow::CollectiveParams const*, std::__cxx11::basic_string<char, std::char_traits<char>, std::allocator<char> > const&, std::function<void (tensorflow::Status const&)>)::{lambda()#3}::operator()() base_collective_executor.cc:0 4 0x0000000000b99403 tensorflow::UnboundedWorkQueue::PooledThreadFunc() ???:0 5 0x0000000000b9f6b1 tensorflow::(anonymous namespace)::PThread::ThreadFn() env.cc:0 6 0x0000000000007ea5 start_thread() pthread_create.c:0 7 0x00000000000feb0d __clone() ???:0 ================================= [g03:35503] *** Process received signal *** [g03:35503] Signal: Segmentation fault (11) [g03:35503] Signal code: (-6) [g03:35503] Failing at address: 0x2ecf700008aaf [g03:35503] [ 0] /lib64/libpthread.so.0(+0xf630)[0x2aaaab7e6630] [g03:35503] [ 1] /usr/ebuild/software/TensorFlow/2.6.0-foss-2021a-CUDA-11.3.1/lib/python3.9/site-packages/tensorflow/python/_pywrap_tensorflow_internal.so(+0x382045b)[0x2aaab68fc45b] [g03:35503] [ 2] /usr/ebuild/software/TensorFlow/2.6.0-foss-2021a-CUDA-11.3.1/lib/python3.9/site-packages/tensorflow/python/_pywrap_tensorflow_internal.so(_ZN10tensorflow11NcclReducer3RunESt8functionIFvRKNS_6StatusEEE+0x1ca)[0x2aaab8d7b88a] [g03:35503] [ 3] /usr/ebuild/software/TensorFlow/2.6.0-foss-2021a-CUDA-11.3.1/lib/python3.9/site-packages/tensorflow/python/../libtensorflow_framework.so.2(+0x9086dc)[0x2aaadc7556dc] [g03:35503] [ 4] /usr/ebuild/software/TensorFlow/2.6.0-foss-2021a-CUDA-11.3.1/lib/python3.9/site-packages/tensorflow/python/../libtensorflow_framework.so.2(_ZN10tensorflow18UnboundedWorkQueue16PooledThreadFuncEv+0x1b3)[0x2aaadc9e6403] [g03:35503] [ 5] /usr/ebuild/software/TensorFlow/2.6.0-foss-2021a-CUDA-11.3.1/lib/python3.9/site-packages/tensorflow/python/../libtensorflow_framework.so.2(+0xb9f6b1)[0x2aaadc9ec6b1] [g03:35503] [ 6] /lib64/libpthread.so.0(+0x7ea5)[0x2aaaab7deea5] [g03:35503] [ 7] /lib64/libc.so.6(clone+0x6d)[0x2aaaac468b0d] [g03:35503] *** End of error message *** -------------------------------------------------------------------------- Primary job terminated normally, but 1 process returned a non-zero exit code. Per user-direction, the job has been aborted. -------------------------------------------------------------------------- -------------------------------------------------------------------------- mpirun noticed that process rank 0 with PID 35503 on node g03 exited on signal 11 (Segmentation fault).
训练代码
from getOneHot import getOneHot from mpi4py import MPI comm = MPI.COMM_WORLD rank = comm.Get_rank() # Load in the parameter files from json import load as loadf with open("params.json", 'r') as inFile: params = loadf(inFile) # Get data files and prep them for the generator from tensorflow import distribute as D callbacks = [] devices = getDevices() print(devices) set_tf_config_mpi() strat = D.experimental.MultiWorkerMirroredStrategy( communication=D.experimental.CollectiveCommunication.NCCL) # Create network from sys import argv resume_training = False print(argv) if "resume_latest" in argv: resume_training = True with strat.scope(): # Scheduler if isinstance(params["learning_rate"], str): # Get the string for the importable function lr = params["learning_rate"] from tensorflow.keras.callbacks import LearningRateScheduler # Use a dummy learning rate params["learning_rate"] = 0.1 # model = create_model(**params) # Get the importable function lr = lr.split(".") baseImport = __import__(lr[0], globals(), locals(), [lr[1]], 0) lr = getattr(baseImport, lr[1]) # Make a schedule lr = LearningRateScheduler(lr) callbacks.append(lr) # Resume Model? model_name = None if resume_training: initial_epoch, model_name = getInitialEpochsAndModelName(rank) if model_name is None: initial_epoch=0 model = create_model(**params) resume_training = False else: from tensorflow.keras.models import load_model model = load_model(model_name) # Load data from disk import numpy if "root" in params.keys(): root = params['root'] else: root = "./" if "filename" in params.keys(): filename = params["filename"] else: filename = "dataset_timeseries.csv" restricted = [ 'euc1', 'e1', 'x1', 'y1', 'z1', 'euc2', 'e2', 'x2', 'y2', 'z2', 'euc3', 'e3', 'x3', 'y3', 'z3', ] x, y = getOneHot("{}/{}".format(root, filename), restricted=restricted, **params) # val_x, val_y = getOneHot("{}/{}".format(root, val_filename), restricted=restricted) val_x, val_y = None, None params["gbatch_size"] = params['batch_size'] * len(devices) print("x.shape =", x.shape) print("y.shape =", y.shape) print("epochs =", params['epochs'], type(params['epochs'])) print("batch =", params['batch_size'], type(params['batch_size'])) print("gbatch =", params["gbatch_size"], type(params["gbatch_size"])) # Load data into a distributed dataset from tensorflow.data import Dataset data = Dataset.from_tensor_slices((x, y)) # Create validation set v = params['validation'] if val_x is not None: vrecord = val_x.shape[0] val = Dataset.from_tensor_slices((val_x, val_y)) validation = val # data.take(vrecord) else: vrecord = int(x.shape[0]*v) validation = data.take(vrecord) validation = validation.batch(params["gbatch_size"]) validation = validation.repeat(params['epochs']) # Validation -- need to do kfold one day # This set should NOT be distributed vsteps = vrecord // params["gbatch_size"] if vrecord % params["gbatch_size"] != 0: vsteps += 1 # Shuffle the data during preprocessing or suffer... # Parallel randomness == nightmare # data = data.shuffle(x.shape[0]) # Ordering these two things is very important! # Consider 3 elements, batch size 2 repeat 2 # [1 2 3] -> [[1 2] [3]] -> [[1 2] [3] [1 2] [3]] (correct) batch -> repeat # [1 2 3] -> [1 2 3 1 2 3] -> [[1 2] [3 1] [2 3]] (incorrect) repeat -> batch # data = data.skip(vrecord) data = data.batch(params["gbatch_size"]) data = data.repeat(params['epochs']) records = x.shape[0] # - vrecord steps = records // params["gbatch_size"] if records % params["gbatch_size"]: steps += 1 print("steps =", steps) # Note that if we are resuming that the number of _remaining_ epochs has # changed! # The number of epochs * steps is the numbers of samples to drop print("initial cardinality = ", data.cardinality()) print("initial v cardinality = ", data.cardinality()) data = data.skip(initial_epoch*steps) validation = validation.skip(initial_epoch*vsteps) print("final cardinality = ", data.cardinality()) print("final v cardinality = ", data.cardinality()) # data = strat.experimental_distribute_dataset(data) # Split into validation and training callbacks = createCallbacks(params, callbacks, rank, resume_training) print(callbacks) history = model.fit(data, epochs=params['epochs'], batch_size=params["gbatch_size"], steps_per_epoch=steps, verbose=0, initial_epoch=initial_epoch, validation_data=validation, validation_steps=vsteps, callbacks=callbacks) if rank == 0: model.save("model-final") else: model.save("checkpoints/model-tmp")
Slurm提交脚本
#!/bin/bash #SBATCH --job-name=RNN_A # Job name #SBATCH --mem=30000 # Job memory request #SBATCH --gres=gpu:4 # Number of requested GPU(s) #SBATCH --time=3-23:00:00 # Time limit days-hrs:min:sec #SBATCH --constraint=rtx_2080 # Specific hardware constraint #SBATCH --error=slurm.err # Error file name #SBATCH --output=slurm.out # Output file name #SBATCH --nodes=1 #SBATCH --ntasks-per-node=1 #SBATCH --cpus-per-task=1 #SBATCH --array=1-2%1 module load Anaconda3/2020.07 module load TensorFlow/2.6.0-foss-2021a-CUDA-11.3.1 module load NCCL/2.10.3-GCCcore-10.3.0-CUDA-11.3.1 mpirun python -u main.py resume_latest
内容的提问来源于stack exchange,提问作者Sri Charan
相关产品推荐
相关产品推荐

