diff --git a/ggml/src/ggml-backend.cpp b/ggml/src/ggml-backend.cpp index 4e36909f45e9..c4c37004f0e4 100644 --- a/ggml/src/ggml-backend.cpp +++ b/ggml/src/ggml-backend.cpp @@ -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; @@ -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] @@ -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 @@ -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 ids; std::vector used_ids; + std::vector missing_ids; for (int split_id = 0; split_id < sched->n_splits; split_id++) { struct ggml_backend_sched_split * split = &splits[split_id]; @@ -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]; @@ -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); @@ -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)]; @@ -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)); + 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; @@ -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)) { @@ -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])); 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; @@ -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); @@ -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;