Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
22 changes: 19 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,13 @@ 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.
const char* ptr = (char*)buffer.data_ptr();
for (const auto& iter : _locked_tensors) {
const char* base = (char*)iter.first;
if (base <= ptr && ptr < base + iter.second) { return true; }
Comment thread
sfc-gh-truwase marked this conversation as resolved.
Outdated
}
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
7 changes: 6 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,8 @@ 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)
{
return _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
5 changes: 5 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,11 @@ 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 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
5 changes: 5 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,11 @@ 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 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
16 changes: 8 additions & 8 deletions deepspeed/runtime/swap_tensor/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@

import torch
from deepspeed.utils.logging import logger
from deepspeed.accelerator import get_accelerator

from deepspeed import comm as dist

Expand Down Expand Up @@ -96,8 +95,8 @@ def get_swap_path(self, offset):

class SwapBufferPool(object):

def __init__(self, buffers):
assert all([get_accelerator().is_pinned(buf) for buf in buffers])
def __init__(self, buffers, aio_handle):
assert all([aio_handle.is_pinned(buf) for buf in buffers])
self.buffers = [SwapBuffer(buf) for buf in buffers]
self.current_index = 0

Expand Down Expand Up @@ -180,14 +179,15 @@ def _get_used_buffers(self):

class SwapBufferManager(object):

def __init__(self, num_elems, count, dtype):
def __init__(self, num_elems, count, dtype, aio_handle):
self.num_elems = num_elems
self.count = count
self.dtype = dtype
self.all_buffers = [
get_accelerator().pin_memory(torch.zeros(num_elems, device='cpu', dtype=dtype), align_bytes=0)
for _ in range(count)
]
# Zero-init: new_cpu_locked_tensor is backed by posix_memalign (uninitialized),
# but optimizer state (e.g. Adam momentum/variance) starts from these buffers and
# must be zeroed, matching the prior torch.zeros() allocation.
dtype_example = torch.zeros(1, device='cpu', dtype=dtype)
self.all_buffers = [aio_handle.new_cpu_locked_tensor(num_elems, dtype_example).zero_() for _ in range(count)]
Comment thread
sfc-gh-truwase marked this conversation as resolved.
self.free_buffer_index = [i for i in range(count)]
self.used_buffer_index = {}
self.gigabytes = (self.all_buffers[0].element_size() * num_elems * count) / (1024**3)
Expand Down
Loading
Loading