Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion csrc/aio/py_lib/deepspeed_cpu_op.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
using namespace std;

cpu_op_desc_t::cpu_op_desc_t(
const std::unique_ptr<struct deepspeed_pin_tensor_t>& pinned_tensor_mgr,
const std::shared_ptr<struct deepspeed_pin_tensor_t>& pinned_tensor_mgr,
const bool read_op,
const torch::Tensor& buffer,
const int fd,
Expand Down
4 changes: 2 additions & 2 deletions csrc/aio/py_lib/deepspeed_cpu_op.h
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,9 @@ struct cpu_op_desc_t : io_op_desc_t {
torch::Tensor _cpu_buffer;
bool _use_bounce_buffer;
bool _is_managed_bounce_buffer;
const std::unique_ptr<struct deepspeed_pin_tensor_t>& _pinned_tensor_mgr;
std::shared_ptr<struct deepspeed_pin_tensor_t> _pinned_tensor_mgr;

cpu_op_desc_t(const std::unique_ptr<struct deepspeed_pin_tensor_t>& pinned_tensor_mgr,
cpu_op_desc_t(const std::shared_ptr<struct deepspeed_pin_tensor_t>& pinned_tensor_mgr,
const bool read_op,
const torch::Tensor& buffer,
const int fd,
Expand Down
25 changes: 22 additions & 3 deletions csrc/aio/py_lib/deepspeed_pin_tensor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,12 @@ deepspeed_pin_tensor_t::~deepspeed_pin_tensor_t()
_locked_tensors.clear();
}

std::shared_ptr<deepspeed_pin_tensor_t> deepspeed_pin_tensor_t::shared()
{
static auto mgr = std::make_shared<deepspeed_pin_tensor_t>();
return mgr;
}

torch::Tensor deepspeed_pin_tensor_t::alloc(const int64_t num_elem,
const torch::TensorOptions& options)
{
Expand All @@ -28,7 +34,10 @@ torch::Tensor deepspeed_pin_tensor_t::alloc(const int64_t num_elem,
auto pinned_buffer = ds_page_aligned_alloc(num_bytes, true);
assert(nullptr != pinned_buffer);

_locked_tensors[pinned_buffer] = num_bytes;
{
std::lock_guard<std::mutex> guard(_mutex);
_locked_tensors[pinned_buffer] = num_bytes;
}

return at::from_blob(pinned_buffer, static_cast<int64_t>(num_elem), options);
}
Expand All @@ -41,6 +50,7 @@ torch::Tensor deepspeed_pin_tensor_t::alloc(const int64_t num_elem, const at::Sc

bool deepspeed_pin_tensor_t::free(torch::Tensor& locked_tensor)
{
std::lock_guard<std::mutex> guard(_mutex);
auto addr = locked_tensor.data_ptr();
if (_locked_tensors.find(addr) != _locked_tensors.end()) {
munlock(addr, _locked_tensors[addr]);
Expand All @@ -55,7 +65,16 @@ bool deepspeed_pin_tensor_t::free(torch::Tensor& locked_tensor)
bool deepspeed_pin_tensor_t::is_managed(const torch::Tensor& buffer)
{
if (!buffer.is_cpu()) { return false; }
auto addr = buffer.data_ptr();
if (_locked_tensors.find(addr) != _locked_tensors.end()) { return true; }
std::lock_guard<std::mutex> guard(_mutex);
// Range check (not exact base match) so slices/views of a locked buffer are
// still recognized as pinned, matching torch's is_pinned() semantics. Require
// the buffer's full byte extent to fall within a single locked region; a buffer
// that starts inside a region but ends past it would have an unpinned tail.
const char* ptr = (char*)buffer.data_ptr();
const char* end = ptr + buffer.nbytes();
for (const auto& iter : _locked_tensors) {
const char* base = (char*)iter.first;
if (base <= ptr && end <= base + iter.second) { return true; }
}
return false;
};
7 changes: 7 additions & 0 deletions csrc/aio/py_lib/deepspeed_pin_tensor.h
Original file line number Diff line number Diff line change
Expand Up @@ -12,15 +12,22 @@ Functionality for managing CPU tensors occupying page-locked memory.
*/

#include <map>
#include <memory>
#include <mutex>
#include "deepspeed_py_aio.h"

struct deepspeed_pin_tensor_t {
std::map<void*, int64_t> _locked_tensors;
std::mutex _mutex;

deepspeed_pin_tensor_t() = default;

~deepspeed_pin_tensor_t();

// Process-wide shared manager so that pinned-buffer recognition is consistent
// across every io handle (each handle references this single instance).
static std::shared_ptr<deepspeed_pin_tensor_t> shared();

torch::Tensor alloc(const int64_t num_elem, const at::ScalarType& elem_type);
torch::Tensor alloc(const int64_t num_elem, const torch::TensorOptions& options);

Expand Down
11 changes: 10 additions & 1 deletion csrc/aio/py_lib/deepspeed_py_io_handle.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ deepspeed_io_handle_t::deepspeed_io_handle_t(const int block_size,
_intra_op_parallelism(intra_op_parallelism),
_aio_config(block_size, queue_depth, single_submit, overlap_events, false),
_num_pending_ops(0),
_pinned_tensor_mgr(new deepspeed_pin_tensor_t())
_pinned_tensor_mgr(deepspeed_pin_tensor_t::shared())
{
for (auto i = 0; i < intra_op_parallelism; ++i) {
_thread_contexts.push_back(std::make_shared<deepspeed_aio_thread_t>(i, _aio_config));
Expand Down Expand Up @@ -376,3 +376,12 @@ bool deepspeed_io_handle_t::free_cpu_locked_tensor(torch::Tensor& locked_tensor)
std::lock_guard<std::mutex> lock(_handle_mutex);
return _pinned_tensor_mgr->free(locked_tensor);
}

bool deepspeed_io_handle_t::is_pinned(const torch::Tensor& buffer)
{
// Mirror cpu_op_desc_t's direct-I/O eligibility check: a buffer needs no bounce
// copy if it is torch-pinned or page-locked by the DeepNVMe manager. Reporting
// only is_managed() would misclassify torch-pinned buffers (e.g. ZeRO-3's fp16
// flat CPU memory) as unpinned and force the extra swap-buffer path.
return buffer.is_pinned() || _pinned_tensor_mgr->is_managed(buffer);
}
4 changes: 3 additions & 1 deletion csrc/aio/py_lib/deepspeed_py_io_handle.h
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ struct deepspeed_io_handle_t {
std::vector<std::thread> _threads;
std::mutex _handle_mutex;
int _num_pending_ops;
std::unique_ptr<struct deepspeed_pin_tensor_t> _pinned_tensor_mgr;
std::shared_ptr<struct deepspeed_pin_tensor_t> _pinned_tensor_mgr;

deepspeed_io_handle_t(const int block_size,
const int queue_depth,
Expand Down Expand Up @@ -78,6 +78,8 @@ struct deepspeed_io_handle_t {

bool free_cpu_locked_tensor(torch::Tensor&);

bool is_pinned(const torch::Tensor& buffer);

int wait();

void _stop_threads();
Expand Down
6 changes: 6 additions & 0 deletions csrc/aio/py_lib/py_ds_aio.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,12 @@ PYBIND11_MODULE(TORCH_EXTENSION_NAME, m)
"tensor"_a,
py::call_guard<py::gil_scoped_release>())

.def("is_pinned",
&deepspeed_aio_handle_t::is_pinned,
"Whether the buffer is directly usable for DeepNVMe I/O (torch-pinned or "
"page-locked by the DeepNVMe pinned-tensor manager).",
"buffer"_a)

.def("wait",
&deepspeed_aio_handle_t::wait,
"Wait for (ongoing) asynchronous operations to complete",
Expand Down
6 changes: 6 additions & 0 deletions csrc/gds/py_lib/py_ds_gds.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,12 @@ PYBIND11_MODULE(TORCH_EXTENSION_NAME, m)
"Free pinned CPU tensor.",
"tensor"_a)

.def("is_pinned",
&deepspeed_gds_handle_t::is_pinned,
"Whether the buffer is directly usable for DeepNVMe I/O (torch-pinned or "
"page-locked by the DeepNVMe pinned-tensor manager).",
"buffer"_a)

.def("new_pinned_device_tensor",
&deepspeed_gds_handle_t::new_pinned_device_tensor,
"Allocate pinned device tensor.",
Expand Down
3 changes: 1 addition & 2 deletions deepspeed/runtime/swap_tensor/async_swapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@
from deepspeed import comm as dist
from deepspeed.utils.logging import logger
from deepspeed.runtime.swap_tensor.utils import swap_out_tensors, SwapBuffer
from deepspeed.accelerator import get_accelerator

INVALID_BUFFER_INDEX = -1
ASYNC_SWAPPER_WAIT_TIMER = 'async_swap_gradient_wait'
Expand Down Expand Up @@ -38,7 +37,7 @@ def has_buffers(self):

def add_buffers(self, buffer_list):
assert len(self.all_buffers) == 0
assert all([get_accelerator().is_pinned(buffer) for buffer in buffer_list])
assert all([self.aio_handle.is_pinned(buffer) for buffer in buffer_list])
dtype = buffer_list[0].dtype
assert all([buffer.dtype == dtype for buffer in buffer_list])

Expand Down
25 changes: 15 additions & 10 deletions deepspeed/runtime/swap_tensor/optimizer_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@
from deepspeed.runtime.swap_tensor.utils import swap_in_tensors, swap_out_tensors, \
MIN_AIO_BYTES, AIO_ALIGNED_BYTES, get_sized_buffers
from deepspeed.runtime.swap_tensor.utils import SwapBufferManager, SwapBufferPool
from deepspeed.accelerator import get_accelerator


class FlattenedTensorSwapInfo(object):
Expand Down Expand Up @@ -44,8 +43,9 @@ def set_buffers(self, compute_buffer, swap_buffer):

class OptimizerStateSwapInfo(object):

def __init__(self, parameter, numel, base_folder):
def __init__(self, parameter, numel, base_folder, aio_handle):
self.tensors = []
self.aio_handle = aio_handle
self.param_id = OptimizerSwapper.parameter_id(parameter)
self.swap_folder = base_folder
self.swapped_gradients = {}
Expand Down Expand Up @@ -93,7 +93,7 @@ def get_swap_paths(self):
def get_swap_buffers_and_paths(self, pinned):
swap_buffers = []
swap_paths = []
select_tensors = [t for t in self.tensors if get_accelerator().is_pinned(t.compute_tensor) == pinned]
select_tensors = [t for t in self.tensors if self.aio_handle.is_pinned(t.compute_tensor) == pinned]
for t in select_tensors:
swap_buffers.append(t.swap_tensor if pinned else t.compute_tensor)
swap_paths.append(t.swap_path)
Expand Down Expand Up @@ -128,7 +128,7 @@ def get_swap_gradient_paths(self):
return [grad.path for grad in self.swapped_gradients.values()]

def get_unpinned_state_tensors(self):
return [t.compute_tensor for t in self.tensors if not get_accelerator().is_pinned(t.compute_tensor)]
return [t.compute_tensor for t in self.tensors if not self.aio_handle.is_pinned(t.compute_tensor)]

def read_unswapped_gradients(self, dest_buffer):
num_elem_count = 0
Expand Down Expand Up @@ -162,9 +162,11 @@ class OptimizerSwapper(object):
def parameter_id(param):
return param.ds_id

def __init__(self, swap_config, aio_config, base_folder, optimizer, largest_numel, device, dtype, timers):
def __init__(self, swap_config, aio_config, base_folder, optimizer, largest_numel, device, dtype, timers,
aio_handle):
self.swap_config = swap_config
self.aio_config = aio_config
self.aio_handle = aio_handle

# NVMe swap management
self.swap_params_info = {}
Expand All @@ -184,7 +186,8 @@ def __init__(self, swap_config, aio_config, base_folder, optimizer, largest_nume
self.dtype = dtype
self.swap_buffer_manager = SwapBufferManager(num_elems=self.largest_numel,
count=swap_config.buffer_count,
dtype=dtype)
dtype=dtype,
aio_handle=aio_handle)

# Timers
self.timers = timers
Expand All @@ -197,6 +200,7 @@ def __init__(self, swap_config, aio_config, base_folder, optimizer, largest_nume
'swap_params_info',
'timers',
'timer_names',
'aio_handle',
]

def purge_state(self):
Expand Down Expand Up @@ -272,7 +276,7 @@ def _initialize_from_swapped_fp16_params(self, aio_handle, fp16_partitions_info,
fp16_pinned_buffers, fp32_parameters):
assert len(fp32_parameters) == len(fp16_partitions_info)
assert len(fp32_parameters) == len(fp16_num_elems)
assert all([get_accelerator().is_pinned(buffer) for buffer in fp16_pinned_buffers])
assert all([aio_handle.is_pinned(buffer) for buffer in fp16_pinned_buffers])

fp32_swap_paths = self._get_swap_paths(parameters=fp32_parameters, num_elems=fp16_num_elems)

Expand All @@ -282,8 +286,8 @@ def _initialize_from_swapped_fp16_params(self, aio_handle, fp16_partitions_info,
assert all([numel >= self.largest_numel for numel in fp16_buffer_numel]), \
f"numel of fp16 buffers {fp16_buffer_numel} is too small for initializing fp32 params {self.largest_numel}"

fp32_swap_buffers = SwapBufferPool(fp32_pinned_buffers)
fp16_swap_buffers = SwapBufferPool(fp16_pinned_buffers)
fp32_swap_buffers = SwapBufferPool(fp32_pinned_buffers, aio_handle=aio_handle)
fp16_swap_buffers = SwapBufferPool(fp16_pinned_buffers, aio_handle=aio_handle)

curr_index = 0
while curr_index < len(fp32_parameters):
Expand Down Expand Up @@ -495,7 +499,8 @@ def _create_param_swap_info(self, parameter, numel):

self.swap_params_info[param_id] = OptimizerStateSwapInfo(parameter=parameter,
numel=numel,
base_folder=self.swap_folder)
base_folder=self.swap_folder,
aio_handle=self.aio_handle)
swap_info = self.swap_params_info[param_id]

self._update_param_state_info(swap_info, parameter)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@
get_sized_buffers
from deepspeed.runtime.swap_tensor.async_swapper import AsyncTensorSwapper
from deepspeed.runtime.swap_tensor.optimizer_utils import OptimizerSwapper
from deepspeed.accelerator import get_accelerator

DEBUG_MODE = False

Expand All @@ -27,16 +26,16 @@
class PartitionedOptimizerSwapper(OptimizerSwapper):

def __init__(self, swap_config, aio_config, base_folder, optimizer, largest_numel, device, dtype, timers):
super(PartitionedOptimizerSwapper, self).__init__(swap_config, aio_config, base_folder, optimizer,
largest_numel, device, dtype, timers)

aio_op = AsyncIOBuilder().load()
self.aio_handle = aio_op.aio_handle(block_size=aio_config[AIO_BLOCK_SIZE],
queue_depth=aio_config[AIO_QUEUE_DEPTH],
single_submit=aio_config[AIO_SINGLE_SUBMIT],
overlap_events=aio_config[AIO_OVERLAP_EVENTS],
intra_op_parallelism=aio_config[AIO_INTRA_OP_PARALLELISM])

super(PartitionedOptimizerSwapper, self).__init__(swap_config, aio_config, base_folder, optimizer,
largest_numel, device, dtype, timers, self.aio_handle)

# Overlap swapping out
self.gradient_swapper = AsyncTensorSwapper(aio_handle=self.aio_handle,
numel_alignment=self.numel_alignment,
Expand Down Expand Up @@ -217,7 +216,7 @@ def _swap_in_gradients(self, aio_handle, parameter, dest_buffer):
if not (swap_info and swap_info.has_gradients()):
return

assert get_accelerator().is_pinned(dest_buffer)
assert self.aio_handle.is_pinned(dest_buffer)
assert parameter.numel() <= dest_buffer.numel()

parameter.grad = dest_buffer.narrow(0, 0, parameter.numel())
Expand Down
13 changes: 6 additions & 7 deletions deepspeed/runtime/swap_tensor/partitioned_param_swapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,7 @@ def _configure_aio(self, ds_config):
if self.use_gds:
self.aio_read_handle.pin_device_tensor(self.buffers)
else:
self.buffers = get_accelerator().pin_memory(self.buffers, align_bytes=0)
self.buffers = self.aio_read_handle.new_cpu_locked_tensor(self.buffers.numel(), self.buffers)

self.swap_out_params = []

Expand Down Expand Up @@ -326,7 +326,7 @@ def swap_in(self, params, async_op=True, swap_in_buffers=None):
def swap_into_buffer(self, param, dest_buffer):
assert param.ds_tensor.status == PartitionedParamStatus.NOT_AVAILABLE, f"param {param.ds_id} is already available or inflight"

require_swap_buffer = not (get_accelerator().is_pinned(dest_buffer)
require_swap_buffer = not (self.aio_read_handle.is_pinned(dest_buffer)
and self._is_io_aligned(dest_buffer.numel()))
Comment thread
sfc-gh-truwase marked this conversation as resolved.

if require_swap_buffer:
Expand Down Expand Up @@ -392,11 +392,10 @@ def _is_io_aligned(self, numel):

def reserve_partitioned_swap_space(self, partition_num_elems):
aligned_numel = sum([self._io_aligned_numel(numel) for numel in partition_num_elems])
self.partitioned_swap_buffer = get_accelerator().pin_memory(torch.zeros(aligned_numel,
device='cpu',
dtype=self.dtype),
align_bytes=0)
self.partitioned_swap_pool = SwapBufferPool([self.partitioned_swap_buffer])
dtype_example = torch.zeros(1, device='cpu', dtype=self.dtype)
self.partitioned_swap_buffer = self.aio_write_handle.new_cpu_locked_tensor(aligned_numel,
dtype_example).zero_()
self.partitioned_swap_pool = SwapBufferPool([self.partitioned_swap_buffer], aio_handle=self.aio_write_handle)

def swap_out_partitioned_params(self, dst_fp16_params, src_fp32_params):
assert self.partitioned_swap_buffer is not None, 'partitioned swap buffers for fp16 params not initialized'
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,9 +52,6 @@ def wait(self):
class PipelinedOptimizerSwapper(OptimizerSwapper):

def __init__(self, swap_config, aio_config, base_folder, optimizer, largest_numel, device, dtype, timers):
super(PipelinedOptimizerSwapper, self).__init__(swap_config, aio_config, base_folder, optimizer, largest_numel,
device, dtype, timers)

aio_op = AsyncIOBuilder().load()
self.write_aio_handle = aio_op.aio_handle(block_size=aio_config[AIO_BLOCK_SIZE],
queue_depth=aio_config[AIO_QUEUE_DEPTH],
Expand All @@ -68,6 +65,9 @@ def __init__(self, swap_config, aio_config, base_folder, optimizer, largest_nume
overlap_events=aio_config[AIO_OVERLAP_EVENTS],
intra_op_parallelism=aio_config[AIO_INTRA_OP_PARALLELISM])

super(PipelinedOptimizerSwapper, self).__init__(swap_config, aio_config, base_folder, optimizer, largest_numel,
device, dtype, timers, self.write_aio_handle)

# Overlap gradient swap out
self.gradient_swapper = AsyncTensorSwapper(aio_handle=self.write_aio_handle,
numel_alignment=self.numel_alignment,
Expand Down
Loading
Loading