Skip to content
Closed
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
125 changes: 97 additions & 28 deletions ggml/src/ggml-backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -771,6 +771,15 @@ struct ggml_backend_sched_split {
struct ggml_cgraph graph;
};

struct ggml_backend_sched_moe_loaded {
ggml_bitset_t * ids;
size_t id_size;
int64_t n_expert;
size_t expert_size;
const void * src_data;
void * dst_data;
};

struct ggml_backend_sched {
bool is_reset; // true if the scheduler has been reset since the last graph split
bool is_alloc;
Expand All @@ -785,6 +794,7 @@ struct ggml_backend_sched {
struct ggml_hash_set hash_set;
int * hv_tensor_backend_ids; // [hash_set.size]
struct ggml_tensor ** hv_tensor_copies; // [hash_set.size][n_backends][n_copies]
struct ggml_backend_sched_moe_loaded * hv_tensor_moe_loaded; // [hash_set.size][n_backends][n_copies]

int * node_backend_ids; // [graph_size]
int * leaf_backend_ids; // [graph_size]
Expand Down Expand Up @@ -829,7 +839,9 @@ struct ggml_backend_sched {

#define hash_id(tensor) ggml_hash_find_or_insert(&sched->hash_set, tensor)
#define tensor_backend_id(tensor) sched->hv_tensor_backend_ids[hash_id(tensor)]
#define tensor_id_copy(id, backend_id, copy_id) sched->hv_tensor_copies[(id) * sched->n_backends * sched->n_copies + (backend_id) * sched->n_copies + (copy_id)]
#define tensor_copy_index(id, backend_id, copy_id) \
((id) * (size_t) sched->n_backends * sched->n_copies + (backend_id) * (size_t) sched->n_copies + (copy_id))
#define tensor_id_copy(id, backend_id, copy_id) sched->hv_tensor_copies[tensor_copy_index(id, backend_id, copy_id)]
#define tensor_copy(tensor, backend_id, copy_id) tensor_id_copy(hash_id(tensor), backend_id, copy_id)

// returns the priority of the backend, lower id is higher priority
Expand Down Expand Up @@ -1545,6 +1557,7 @@ static enum ggml_status ggml_backend_sched_compute_splits(ggml_backend_sched_t s
ggml_tensor * prev_ids_tensor = nullptr;
std::vector<int32_t> ids;
std::vector<ggml_bitset_t> used_ids;
std::vector<ggml_bitset_t> missing_ids;

for (int split_id = 0; split_id < sched->n_splits; split_id++) {
struct ggml_backend_sched_split * split = &splits[split_id];
Expand All @@ -1566,12 +1579,14 @@ static enum ggml_status ggml_backend_sched_compute_splits(ggml_backend_sched_t s
}
ggml_backend_tensor_copy(input, input_cpy);
} else {
// wait for the split backend to finish using the input before overwriting it
if (sched->events[split_backend_id][sched->cur_copy] != NULL) {
ggml_backend_event_wait(split_backend, sched->events[split_backend_id][sched->cur_copy]);
} else {
ggml_backend_synchronize(split_backend);
}
auto wait_for_split_input = [&]() {
// wait for the split backend to finish using the input before overwriting it
if (sched->events[split_backend_id][sched->cur_copy] != NULL) {
ggml_backend_event_wait(split_backend, sched->events[split_backend_id][sched->cur_copy]);
} else {
ggml_backend_synchronize(split_backend);
}
};

// when offloading MoE weights, we can reduce the amount of data copied by copying only the experts that are used
ggml_tensor * node = split->graph.nodes[0];
Expand All @@ -1581,9 +1596,10 @@ static enum ggml_status ggml_backend_sched_compute_splits(ggml_backend_sched_t s
(node->src[0] == input_cpy && node->op == GGML_OP_MUL_MAT_ID)
//|| (node->src[1] == input_cpy && node->op == GGML_OP_ADD_ID) /* GGML_OP_ADD_ID weights are small and not worth splitting */
)) {

const int64_t n_expert = node->op == GGML_OP_MUL_MAT_ID ? input->ne[2] : input->ne[1];
const size_t expert_size = node->op == GGML_OP_MUL_MAT_ID ? input->nb[2] : input->nb[1];
GGML_ASSERT(n_expert > 0);
const size_t expert_id_size = ggml_bitset_size(n_expert);

ggml_backend_synchronize(input_backend);

Expand All @@ -1601,14 +1617,14 @@ static enum ggml_status ggml_backend_sched_compute_splits(ggml_backend_sched_t s
}
}

if (ids_tensor != prev_ids_tensor) {
if (ids_tensor != prev_ids_tensor || used_ids.size() != expert_id_size) {
ids.resize(ggml_nbytes(ids_tensor) / sizeof(int32_t));
ggml_backend_tensor_get_async(ids_backend, ids_tensor, ids.data(), 0, ggml_nbytes(ids_tensor));
ggml_backend_synchronize(ids_backend);

// find the used experts
used_ids.clear();
used_ids.resize(ggml_bitset_size(n_expert));
used_ids.resize(expert_id_size);
for (int64_t i1 = 0; i1 < ids_tensor->ne[1]; i1++) {
for (int64_t i0 = 0; i0 < ids_tensor->ne[0]; i0++) {
int32_t id = ids[i1 * ids_tensor->nb[1]/sizeof(int32_t) + i0 * ids_tensor->nb[0]/sizeof(int32_t)];
Expand All @@ -1620,6 +1636,38 @@ static enum ggml_status ggml_backend_sched_compute_splits(ggml_backend_sched_t s
prev_ids_tensor = ids_tensor;
}

const size_t input_hash_id = hash_id(input);
const size_t copy_idx = tensor_copy_index(input_hash_id, split_backend_id, sched->cur_copy);
const size_t loaded_id_size = expert_id_size;
ggml_backend_sched_moe_loaded * loaded = &sched->hv_tensor_moe_loaded[copy_idx];

if (loaded->id_size != loaded_id_size) {
free(loaded->ids);
loaded->ids = (ggml_bitset_t *) calloc(loaded_id_size, sizeof(ggml_bitset_t));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The memory allocation for loaded->ids using calloc should be checked for failure, especially since expert_id_size could be large depending on the model configuration.

loaded->id_size = loaded_id_size;
loaded->n_expert = 0;
loaded->expert_size = 0;
}

ggml_bitset_t * loaded_ids = loaded->ids;
GGML_ASSERT(loaded_ids != nullptr);

if (loaded->n_expert != n_expert || loaded->expert_size != expert_size ||
loaded->src_data != input->data || loaded->dst_data != input_cpy->data) {
memset(loaded_ids, 0, loaded_id_size * sizeof(ggml_bitset_t));
loaded->n_expert = n_expert;
loaded->expert_size = expert_size;
loaded->src_data = input->data;
loaded->dst_data = input_cpy->data;
}

bool has_missing_ids = false;
missing_ids.resize(loaded_id_size);
for (size_t i = 0; i < loaded_id_size; ++i) {
missing_ids[i] = used_ids[i] & ~loaded_ids[i];
has_missing_ids = has_missing_ids || missing_ids[i] != 0;
}

// group consecutive experts and copy them together
auto copy_experts = [&](int32_t first_id, int32_t last_id) {
const size_t expert_offset = first_id * expert_size;
Expand All @@ -1633,32 +1681,43 @@ static enum ggml_status ggml_backend_sched_compute_splits(ggml_backend_sched_t s
// copy a bit extra at the to ensure there are no NaNs in the padding of the last expert
// this is necessary for MMQ in the CUDA backend
expert_size_copy + padding_end);

for (int32_t id = first_id; id <= last_id; ++id) {
ggml_bitset_set(loaded_ids, id);
}
};

int id = 0;
while (!ggml_bitset_get(used_ids.data(), id)) {
id++;
}
int32_t first_id = id;
int32_t last_id = first_id;
if (has_missing_ids) {
wait_for_split_input();

for (++id; id < n_expert; ++id) {
if (!ggml_bitset_get(used_ids.data(), id)) {
continue;
int id = 0;
while (id < n_expert && !ggml_bitset_get(missing_ids.data(), id)) {
id++;
}
GGML_ASSERT(id < n_expert);
int32_t first_id = id;
int32_t last_id = first_id;

for (++id; id < n_expert; ++id) {
if (!ggml_bitset_get(missing_ids.data(), id)) {
continue;
}

if (id == last_id + 1) {
last_id = id;
continue;
}

if (id == last_id + 1) {
copy_experts(first_id, last_id);

first_id = id;
last_id = id;
continue;
}

copy_experts(first_id, last_id);

first_id = id;
last_id = id;
}
copy_experts(first_id, last_id);
} else {
wait_for_split_input();

// try async copy, but if not possible, we can still use a sync copy without synchronizing the dst backend, since we handle the synchronization here with multiple copies and events
// TODO: add public function to facilitate this, since applications do not have direct access to the backend interface
if (!split_backend->iface.cpy_tensor_async || !split_backend->iface.cpy_tensor_async(input_backend, split_backend, input, input_cpy)) {
Expand Down Expand Up @@ -1754,7 +1813,9 @@ ggml_backend_sched_t ggml_backend_sched_new(
// FIXME: needs to be size*2 to account for leafs (do it in graph_split instead)
sched->hash_set = ggml_hash_set_new(graph_size);
sched->hv_tensor_backend_ids = (int *) malloc(sched->hash_set.size * sizeof(sched->hv_tensor_backend_ids[0]));
sched->hv_tensor_copies = (ggml_tensor **) malloc(sched->hash_set.size * sched->n_backends * sched->n_copies * sizeof(struct ggml_tensor *));
const size_t tensor_copy_count = sched->hash_set.size * sched->n_backends * sched->n_copies;
sched->hv_tensor_copies = (ggml_tensor **) malloc(tensor_copy_count * sizeof(struct ggml_tensor *));
sched->hv_tensor_moe_loaded = (ggml_backend_sched_moe_loaded *) calloc(tensor_copy_count, sizeof(sched->hv_tensor_moe_loaded[0]));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The allocation of hv_tensor_moe_loaded using calloc should be checked for NULL to ensure the scheduler initialization is robust against memory pressure.


const size_t ggml_sched_max_splits = graph_size; // at most there is one split for each node in the graph
const size_t nodes_size = graph_size + ggml_sched_max_splits*GGML_SCHED_MAX_SPLIT_INPUTS*2;
Expand Down Expand Up @@ -1804,10 +1865,17 @@ void ggml_backend_sched_free(ggml_backend_sched_t sched) {
}
ggml_gallocr_free(sched->galloc);
ggml_free(sched->ctx);
const size_t tensor_copy_count = sched->hash_set.size * sched->n_backends * sched->n_copies;
if (sched->hv_tensor_moe_loaded != NULL) {
for (size_t i = 0; i < tensor_copy_count; ++i) {
free(sched->hv_tensor_moe_loaded[i].ids);
}
}
ggml_hash_set_free(&sched->hash_set);
free(sched->splits);
free(sched->hv_tensor_backend_ids);
free(sched->hv_tensor_copies);
free(sched->hv_tensor_moe_loaded);
free(sched->node_backend_ids);
free(sched->leaf_backend_ids);
free(sched->prev_node_backend_ids);
Expand All @@ -1824,7 +1892,8 @@ void ggml_backend_sched_reset(ggml_backend_sched_t sched) {
if (!sched->is_reset) {
ggml_hash_set_reset(&sched->hash_set);
memset(sched->hv_tensor_backend_ids, -1, sched->hash_set.size * sizeof(sched->hv_tensor_backend_ids[0]));
memset(sched->hv_tensor_copies, 0, sched->hash_set.size * sched->n_backends * sched->n_copies * sizeof(struct ggml_tensor *));
const size_t tensor_copy_count = sched->hash_set.size * sched->n_backends * sched->n_copies;
memset(sched->hv_tensor_copies, 0, tensor_copy_count * sizeof(struct ggml_tensor *));
sched->is_reset = true;
}
sched->is_alloc = false;
Expand Down