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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 14:40:16