Download cpp/sdq_pager.cpp from HelloSun/sddqwen35a3b: direct link, hf CLI and curl.
- Browser
- Download file 74.6 kB
-
https://huggingface.co/HelloSun/sddqwen35a3b/resolve/main/cpp/sdq_pager.cpp
- Command line
-
hf download hf://HelloSun/sddqwen35a3b/cpp/sdq_pager.cpp
-
curl -L -o sdq_pager.cpp https://huggingface.co/HelloSun/sddqwen35a3b/resolve/main/cpp/sdq_pager.cpp
74.6 kB
| // SDQ expert pager 實作:見 sdq_pager.h 的設計說明 | |
| namespace sdq { | |
| // ============================================================ 小工具 | |
| uint64_t now_us() { | |
| struct timespec ts; | |
| clock_gettime(CLOCK_MONOTONIC, &ts); | |
| return (uint64_t) ts.tv_sec * 1000000ull + (uint64_t) ts.tv_nsec / 1000ull; | |
| } | |
| size_t process_rss_bytes() { | |
| std::ifstream f("/proc/self/statm"); | |
| size_t total = 0, resident = 0; | |
| if (f >> total >> resident) { | |
| return resident * (size_t) sysconf(_SC_PAGESIZE); | |
| } | |
| return 0; | |
| } | |
| size_t physical_ram_bytes() { | |
| std::ifstream f("/proc/meminfo"); | |
| std::string key; | |
| size_t value = 0; | |
| std::string unit; | |
| while (f >> key >> value >> unit) { | |
| if (key == "MemTotal:") { | |
| return value * 1024; | |
| } | |
| } | |
| return 0; | |
| } | |
| static size_t env_size(const char * name, size_t dflt) { | |
| const char * v = getenv(name); | |
| if (!v || !*v) { | |
| return dflt; | |
| } | |
| return (size_t) strtoull(v, nullptr, 10); | |
| } | |
| static bool env_flag(const char * name, bool dflt) { | |
| const char * v = getenv(name); | |
| if (!v || !*v) { | |
| return dflt; | |
| } | |
| return !(v[0] == '0' && v[1] == '\0'); | |
| } | |
| // ============================================================ GGUF 檔頭解析 | |
| namespace { | |
| struct gguf_reader { | |
| const uint8_t * p = nullptr; | |
| size_t size = 0; | |
| size_t o = 0; | |
| uint8_t u8() { return p[o++]; } | |
| uint32_t u32() { uint32_t v; memcpy(&v, p + o, 4); o += 4; return v; } | |
| uint64_t u64() { uint64_t v; memcpy(&v, p + o, 8); o += 8; return v; } | |
| int64_t i64() { int64_t v; memcpy(&v, p + o, 8); o += 8; return v; } | |
| float f32() { float v; memcpy(&v, p + o, 4); o += 4; return v; } | |
| std::string str() { | |
| uint64_t n = u64(); | |
| std::string s((const char *) p + o, (size_t) n); | |
| o += n; | |
| return s; | |
| } | |
| void skip(uint64_t n) { o += (size_t) n; } | |
| }; | |
| struct kv_value { | |
| enum kind { NONE, INT, REAL, STR, ARR } k = NONE; | |
| int64_t i = 0; | |
| double f = 0; | |
| std::string s; | |
| std::vector<int64_t> arr; | |
| }; | |
| struct tensor_info { | |
| std::string name; | |
| std::vector<int64_t> dims; | |
| ggml_type type = GGML_TYPE_COUNT; | |
| uint64_t offset = 0; | |
| }; | |
| size_t ggml_type_size_of(ggml_type t) { | |
| return ggml_row_size(t, 1); | |
| } | |
| } // namespace | |
| // ============================================================ pager | |
| pager & pager::instance() { | |
| static pager p; | |
| return p; | |
| } | |
| void pager::load_gguf_header(const std::string & path) { | |
| const int fd = open(path.c_str(), O_RDONLY); | |
| if (fd < 0) { | |
| throw std::runtime_error("無法開啟 GGUF:" + path); | |
| } | |
| struct stat st; | |
| fstat(fd, &st); | |
| // 只讀檔頭:metadata + tensor 資訊大約 11 MB | |
| const size_t want = (size_t) std::min<int64_t>(st.st_size, 64ull << 20); | |
| std::vector<uint8_t> buf(want); | |
| ssize_t got = 0; | |
| while (got < (ssize_t) want) { | |
| ssize_t n = pread(fd, buf.data() + got, want - got, got); | |
| if (n <= 0) { | |
| break; | |
| } | |
| got += n; | |
| } | |
| buf.resize((size_t) std::max<ssize_t>(got, 0)); | |
| gguf_reader r{buf.data(), buf.size(), 0}; | |
| if (memcmp(buf.data(), "GGUF", 4) != 0) { | |
| close(fd); | |
| throw std::runtime_error("不是 GGUF 檔:" + path); | |
| } | |
| r.o = 4; | |
| const uint32_t version = r.u32(); | |
| if (version < 2 || version > 3) { | |
| close(fd); | |
| throw std::runtime_error("不支援的 GGUF 版本"); | |
| } | |
| const uint64_t n_tensors = r.u64(); | |
| const uint64_t n_kv = r.u64(); | |
| std::map<std::string, kv_value> kv; | |
| for (uint64_t i = 0; i < n_kv; i++) { | |
| std::string key = r.str(); | |
| const uint32_t t = r.u32(); | |
| kv_value v; | |
| switch (t) { | |
| case 0: v.k = kv_value::INT; v.i = (int64_t) r.u32(); v.f = (double) v.i; break; | |
| case 1: v.k = kv_value::INT; v.i = (int64_t) (int8_t) r.u32(); v.f = (double) v.i; break; | |
| case 2: v.k = kv_value::INT; v.i = r.u32(); v.f = (double) v.i; break; | |
| case 3: v.k = kv_value::INT; v.i = (int16_t) r.u32(); v.f = (double) v.i; break; | |
| case 4: v.k = kv_value::INT; v.i = (int64_t) (uint32_t) r.u32(); v.f = (double) v.i; break; | |
| case 5: v.k = kv_value::INT; v.i = (int32_t) r.u32(); v.f = (double) v.i; break; | |
| case 6: v.k = kv_value::REAL; v.f = r.f32(); v.i = (int64_t) v.f; break; | |
| case 7: v.k = kv_value::INT; v.i = r.u8(); v.f = (double) v.i; break; // BOOL = 1 byte | |
| case 8: v.k = kv_value::STR; v.s = r.str(); break; | |
| case 9: { | |
| v.k = kv_value::ARR; | |
| const uint32_t et = r.u32(); | |
| const uint64_t n = r.u64(); | |
| for (uint64_t j = 0; j < n; j++) { | |
| switch (et) { | |
| case 2: v.arr.push_back(r.u32()); break; | |
| case 3: v.arr.push_back((int16_t) r.u32()); break; | |
| case 4: v.arr.push_back((uint32_t) r.u32()); break; | |
| case 5: v.arr.push_back((int32_t) r.u32()); break; | |
| case 6: v.arr.push_back((int64_t) r.f32()); break; | |
| case 7: v.arr.push_back(r.u32()); break; | |
| case 10: v.arr.push_back((int64_t) r.u64()); break; | |
| case 11: v.arr.push_back(r.i64()); break; | |
| case 12: v.arr.push_back((int64_t) (double) 0); break; // 極少見 | |
| case 8: r.skip(r.u64()); v.arr.push_back(0); break; | |
| default: break; | |
| } | |
| } | |
| } break; | |
| case 10: v.k = kv_value::INT; v.i = (int64_t) r.u64(); v.f = (double) v.i; break; | |
| case 11: v.k = kv_value::INT; v.i = r.i64(); v.f = (double) v.i; break; | |
| case 12: v.k = kv_value::REAL; { double d; memcpy(&d, buf.data() + r.o, 8); r.o += 8; v.f = d; v.i = (int64_t) d; } break; | |
| default: close(fd); throw std::runtime_error("GGUF kv 型別不支援: " + std::to_string(t)); | |
| } | |
| kv[key] = v; | |
| } | |
| std::vector<tensor_info> tensors; | |
| tensors.reserve(n_tensors); | |
| for (uint64_t i = 0; i < n_tensors; i++) { | |
| tensor_info ti; | |
| ti.name = r.str(); | |
| const uint32_t nd = r.u32(); | |
| for (uint32_t d = 0; d < nd; d++) { | |
| ti.dims.push_back((int64_t) r.u64()); | |
| } | |
| ti.type = (ggml_type) r.u32(); | |
| ti.offset = r.u64(); | |
| tensors.push_back(std::move(ti)); | |
| } | |
| const uint64_t align = kv.count("general.alignment") ? (uint64_t) kv["general.alignment"].i : 32; | |
| const uint64_t data_start = (r.o + align - 1) / align * align; | |
| close(fd); | |
| // ---- hparams | |
| auto gi = [&](const std::string & k, int64_t dflt) -> int64_t { | |
| auto it = kv.find(k); | |
| return it == kv.end() ? dflt : it->second.i; | |
| }; | |
| info_.arch = kv.count("general.architecture") ? kv["general.architecture"].s : ""; | |
| info_.name = kv.count("general.name") ? kv["general.name"].s : ""; | |
| const std::string pfx = info_.arch + "."; | |
| info_.n_layer = gi(pfx + "block_count", 0); | |
| info_.n_expert = gi(pfx + "expert_count", 0); | |
| info_.n_expert_used= gi(pfx + "expert_used_count", 0); | |
| info_.n_embd = gi(pfx + "embedding_length", 0); | |
| info_.n_ff_exp = gi(pfx + "expert_feed_forward_length", 0); | |
| info_.n_ff_shexp = gi(pfx + "expert_shared_feed_forward_length", 0); | |
| info_.n_head = gi(pfx + "attention.head_count", 0); | |
| info_.n_head_kv = gi(pfx + "attention.head_count_kv", 0); | |
| info_.n_embd_head_k= gi(pfx + "attention.key_length", 0); | |
| info_.n_embd_head_v= gi(pfx + "attention.value_length", info_.n_embd_head_k); | |
| info_.context_length = gi(pfx + "context_length", 0); | |
| info_.full_attention_interval = gi(pfx + "full_attention_interval", 4); | |
| if (info_.n_expert <= 0 || info_.n_layer <= 0) { | |
| throw std::runtime_error("這不是 MoE 模型(expert_count=" + std::to_string(info_.n_expert) + ")"); | |
| } | |
| // ---- 頁面表:每個 (layer, expert) 的三段 | |
| layers_.assign((size_t) info_.n_layer, {}); | |
| std::map<std::string, const tensor_info *> by_name; | |
| for (const auto & ti : tensors) { | |
| by_name[ti.name] = &ti; | |
| } | |
| size_t max_expert_bytes = 0; | |
| for (int64_t il = 0; il < info_.n_layer; il++) { | |
| const std::string b = "blk." + std::to_string(il) + "."; | |
| layer_desc & L = layers_[(size_t) il]; | |
| L.experts.assign((size_t) info_.n_expert, {}); | |
| auto get = [&](const std::string & suffix) -> const tensor_info * { | |
| auto it = by_name.find(b + suffix); | |
| return it == by_name.end() ? nullptr : it->second; | |
| }; | |
| const tensor_info * t_gate = get("ffn_gate_exps.weight"); | |
| const tensor_info * t_up = get("ffn_up_exps.weight"); | |
| const tensor_info * t_gu = get("ffn_gate_up_exps.weight"); | |
| const tensor_info * t_down = get("ffn_down_exps.weight"); | |
| if (!t_down || (t_gate && !t_up && !t_gu)) { | |
| continue; // 這層不是 MoE(理論上不會發生) | |
| } | |
| for (int64_t e = 0; e < info_.n_expert; e++) { | |
| expert_desc & ed = L.experts[(size_t) e]; | |
| if (t_gate || t_gu) { | |
| const tensor_info * t = t_gate ? t_gate : t_gu; | |
| const uint64_t nb = (uint64_t) ggml_row_size(t->type, t->dims[0]) * (uint64_t) t->dims[1]; | |
| ed.seg[part::GATE].off = data_start + t->offset + (uint64_t) e * nb; | |
| ed.seg[part::GATE].size = nb; | |
| ed.ne0 = t->dims[0]; | |
| ed.ne1 = t->dims[1]; | |
| ed.type[part::GATE] = t->type; | |
| } | |
| if (t_up) { | |
| const uint64_t nb = (uint64_t) ggml_row_size(t_up->type, t_up->dims[0]) * (uint64_t) t_up->dims[1]; | |
| ed.seg[part::UP].off = data_start + t_up->offset + (uint64_t) e * nb; | |
| ed.seg[part::UP].size = nb; | |
| ed.type[part::UP] = t_up->type; // gate/up 都是 Q4_K,型別一定要各自填 | |
| if (!t_gate) { | |
| ed.ne0 = t_up->dims[0]; | |
| ed.ne1 = t_up->dims[1]; | |
| ed.type[part::GATE] = t_up->type; | |
| } | |
| } | |
| { | |
| const uint64_t nb = (uint64_t) ggml_row_size(t_down->type, t_down->dims[0]) * (uint64_t) t_down->dims[1]; | |
| ed.seg[part::DOWN].off = data_start + t_down->offset + (uint64_t) e * nb; | |
| ed.seg[part::DOWN].size = nb; | |
| ed.type[part::DOWN] = t_down->type; | |
| ed.down_ne0 = t_down->dims[0]; | |
| ed.down_ne1 = t_down->dims[1]; | |
| } | |
| ed.fused_gate_up = (t_gate == nullptr && t_gu != nullptr); | |
| ed.bytes = 0; | |
| int64_t acc = 0; | |
| for (int p = 0; p < 3; p++) { | |
| if (ed.seg[p].size == 0) { | |
| continue; | |
| } | |
| acc = (acc + 63) & ~(int64_t) 63; // 每段 64B 對齊 | |
| ed.slot_off[p] = acc; | |
| acc += (int64_t) ed.seg[p].size; | |
| } | |
| ed.bytes = acc; | |
| max_expert_bytes = std::max(max_expert_bytes, (size_t) ed.bytes); | |
| } | |
| L.bytes_per_expert = L.experts.empty() ? 0 : L.experts[0].bytes; | |
| } | |
| slot_bytes_ = (max_expert_bytes + 4095) & ~(size_t) 4095; | |
| if (slot_bytes_ < (256ull << 10)) { | |
| slot_bytes_ = 256ull << 10; | |
| } | |
| // ---- 大小兩種槽位 ---------------------------------------------------- | |
| // 這個模型有 37 層的 expert 是 1900544 B(Q4_K + Q5_K),但層 34/38/39 是 | |
| // 2039808 B(down 換成 Q6_K)。若所有槽位都照最大的配,arena 有 6.8% 被浪費 | |
| // ——而 arena 就是 8 GB 預算裡最寶貴的資源(實測 +5% 槽位 ≈ +9% tok/s)。 | |
| // | |
| // 解法:槽位分成兩區。「小槽位」給 37 層用,「大槽位」只給那 3 層用, | |
| // 大槽位的數量按「那 3 層佔請求的比例」分配(3/40 = 7.5%)。 | |
| // 選「大多數層」當小的:用眾數而不是最大值,這樣即使某個模型只有一種大小, | |
| // 也不會誤判。 | |
| { | |
| std::map<size_t, int> hist; | |
| for (const auto & L : layers_) { | |
| hist[(size_t) L.bytes_per_expert]++; | |
| } | |
| size_t small = 0; | |
| int best_n = -1; | |
| for (const auto & kv : hist) { | |
| if (kv.second > best_n) { | |
| best_n = kv.second; | |
| small = kv.first; | |
| } | |
| } | |
| const size_t mx = (size_t) (layers_.empty() ? 0 : layers_[0].bytes_per_expert); | |
| size_t big = 0; | |
| for (const auto & kv : hist) { | |
| (void) kv; | |
| } | |
| big = 0; | |
| for (const auto & L : layers_) { | |
| if ((size_t) L.bytes_per_expert > small) { | |
| big = std::max(big, (size_t) L.bytes_per_expert); | |
| } | |
| } | |
| small_bytes_ = (small + 4095) & ~(size_t) 4095; | |
| big_bytes_ = big ? ((big + 4095) & ~(size_t) 4095) : 0; | |
| if (big_bytes_ && big_bytes_ <= small_bytes_) { | |
| big_bytes_ = 0; // 其實一樣大,別分兩區 | |
| } | |
| slot_bytes_ = big_bytes_ ? big_bytes_ : small_bytes_; | |
| int n_big_layers = 0; | |
| for (auto & L : layers_) { | |
| L.needs_big_slot = big_bytes_ && (size_t) L.bytes_per_expert > small_bytes_; | |
| if (L.needs_big_slot) { | |
| n_big_layers++; | |
| } | |
| } | |
| (void) mx; | |
| fprintf(stderr, | |
| "[sdq] 槽位大小:小的 %zu B(%zu 層)/大的 %zu B(%d 層)\n", | |
| small_bytes_, (int) layers_.size() - n_big_layers, big_bytes_, n_big_layers); | |
| } | |
| } | |
| bool pager::need_big_slot(int il, int ie) const { | |
| if (!big_bytes_) { | |
| return false; | |
| } | |
| const layer_desc * L = layer(il); | |
| if (!L) { | |
| return false; | |
| } | |
| if (L->needs_big_slot) { | |
| return true; | |
| } | |
| return ie >= 0 && ie < (int) L->experts.size() && | |
| (size_t) L->experts[(size_t) ie].bytes > small_bytes_; | |
| } | |
| void pager::alloc_arena() { | |
| const size_t slot = cfg_.slot_bytes ? cfg_.slot_bytes : slot_bytes_; | |
| slot_bytes_ = slot; | |
| // arena 大小決定:RAM 是「總預算」,不是「還有多少可用」 | |
| // 總預算 = 非專家權重(2.4 GB) + KV cache + ggml 計算緩衝 + 詞表 + arena | |
| // 所以有指定 ram_budget 時,用「目前 RSS(不含 arena)」當作已用額度扣掉。 | |
| // 預留量要涵蓋 arena 之外、但「在 pager 算完之後才長出來」的記憶體: | |
| // llama.cpp 的 CPU compute buffer(125 MiB @ubatch=128、502 MiB @512) | |
| // + KV cache(ctx=4096 約 670 MiB)+ 圖與執行緒堆疊 | |
| // 這些在 sizing 的當下還沒配置,所以不會出現在 rss_now 裡,必須手動扣。 | |
| // 設太小 = 實際用量會超過 RAM 預算。實測(RLIMIT_DATA 8 GB、ctx 4096、681 token 提示詞): | |
| // 預留 700 MB → 峰值 RSS 8231 MB(超出 8 GB 預算 39 MB) | |
| // 預留 900 MB → 峰值 RSS 8035 MB(留 157 MB 餘量) | |
| // 代價是槽位數:2039 → 1936 槽,decode 5.85 → 5.37 tok/s(約 -8%)。 | |
| // 寧可少 8% 速度,也不要在目標機器的 8 GB 上 OOM —— 這是本專案的硬約束。 | |
| const size_t reserve = (size_t) env_size("SDQ_RESERVE_MB", 900) << 20; | |
| const size_t rss_now = process_rss_bytes(); | |
| size_t arena; | |
| if (cfg_.ram_budget_bytes) { | |
| const size_t used = rss_now; // 這時候 arena 還沒配置 | |
| size_t avail = cfg_.ram_budget_bytes > used + reserve | |
| ? cfg_.ram_budget_bytes - used - reserve | |
| : 512ull << 20; | |
| arena = avail; | |
| if (cfg_.arena_bytes) { | |
| arena = std::min(arena, cfg_.arena_bytes); | |
| } | |
| fprintf(stderr, "[sdq] RAM 預算 %.2f GiB:已用(RSS) %.2f GiB、預留 %.2f GiB → arena %.2f GiB\n", | |
| cfg_.ram_budget_bytes / 1073741824.0, used / 1073741824.0, | |
| reserve / 1073741824.0, arena / 1073741824.0); | |
| } else { | |
| const size_t phys = physical_ram_bytes(); | |
| const size_t budget = (size_t) (phys * 0.70); | |
| arena = cfg_.arena_bytes ? cfg_.arena_bytes : std::min(budget, (size_t) (4ull << 30)); | |
| } | |
| if (arena < slot_bytes_ * 16) { | |
| throw std::runtime_error("expert arena 太小"); | |
| } | |
| // 切成兩區:大槽位只分給「需要大槽位的層」,比例 = 那些層佔全部層的比例。 | |
| const bool force_uniform = env_flag("SDQ_FORCE_UNIFORM_SLOT", false); | |
| if (big_bytes_ && small_bytes_ && !layers_.empty() && !force_uniform) { | |
| int n_big_layers = 0; | |
| for (const auto & L : layers_) { | |
| if (L.needs_big_slot) { | |
| n_big_layers++; | |
| } | |
| } | |
| const double f = (double) n_big_layers / (double) layers_.size(); | |
| const double avg = f * (double) big_bytes_ + (1.0 - f) * (double) small_bytes_; | |
| size_t n_total = (size_t)((double) arena / avg); | |
| size_t nb = (size_t)(n_total * f); | |
| // 大槽位的下限不是「有點就好」,而是**工作集的下限**: | |
| // 一個 token 在每個大槽位層要 8 個大 expert,這個模型有 3 個大槽位層 | |
| // → 至少要 3 × 8 = 24 個大槽位,而且還要留給預取的空間。 | |
| // 只給 16 個時量到過「找不到能放 il=34 ie=114 的槽位」而中止 | |
| // (這是在 verify 新增的預取壓力測試 [2b] 抓到的)。 | |
| const size_t big_need = (size_t) (n_big_layers * info_.n_expert_used); | |
| if (nb < big_need) { | |
| nb = big_need; | |
| } | |
| if (nb > n_total / 4) { | |
| nb = n_total / 4; | |
| } | |
| // 連工作集下限都放不進去時,寧可讓小槽位少一點(n_total 保持不變), | |
| // 否則會算出一個「必定會中止」的配置。 | |
| n_big_ = nb; | |
| n_small_ = n_total > nb ? n_total - nb : 0; | |
| arena = n_small_ * small_bytes_ + n_big_ * big_bytes_; | |
| } else { | |
| // 單一尺寸:全部槽位都是「大的那一種」(slot_bytes_)。 | |
| // 這是 SDQ_FORCE_UNIFORM_SLOT=1 與「模型只有一種 expert 大小」時的路徑。 | |
| // | |
| // 這裡必須記成 n_small_=0 / n_big_=全部,不能記成「n_small_=全部、單位用大的」: | |
| // slot_ptr()/slot_size_of() 是用「索引 < n_small_ → 小槽位」來分支的, | |
| // 把它記成全部都是小槽位的話,大 expert 會拿到小槽位的 stride | |
| // (寫出槽位邊界、算錯),而且大 expert 永遠配不到槽位 → 推論卡死。 | |
| // (第一版就是這樣,實驗跑到一半 hang 才找出來。) | |
| const size_t unit = slot_bytes_ ? slot_bytes_ : small_bytes_; | |
| n_small_ = 0; | |
| n_big_ = arena / unit; | |
| arena = n_big_ * unit; | |
| } | |
| if (n_small_ == 0 && n_big_ == 0) { | |
| throw std::runtime_error("expert arena 太小(算不出槽位)"); | |
| } | |
| n_slots_ = n_small_ + n_big_; | |
| if (n_slots_ < 16) { | |
| throw std::runtime_error("expert arena 太小(算不出 16 個槽位)"); | |
| } | |
| void * p = nullptr; | |
| const size_t alloc = arena + (2ull << 20); | |
| if (posix_memalign(&p, 2ull << 20, alloc) != 0 || !p) { | |
| throw std::runtime_error("無法配置 expert arena(" + std::to_string(arena >> 20) + " MB)"); | |
| } | |
| // 讓 arena 全部變成實體 RAM,並提示核心使用 huge page(減少 TLB 壓力) | |
| memset(p, 0, alloc); | |
| madvise(p, alloc, MADV_HUGEPAGE); | |
| arena_ = (uint8_t *) p; | |
| arena_bytes_ = arena; | |
| // OPT-1:剖面計數表(n_layer × n_expert,40×256 = 10240 筆 ≈ 40 KB,可忽略) | |
| prof_.assign((size_t) info_.n_layer * (size_t) info_.n_expert, 0); | |
| if (!cfg_.profile_in.empty()) { | |
| load_profile(cfg_.profile_in); | |
| } | |
| profile_out_ = cfg_.profile_out; | |
| slots_.assign(n_slots_, {}); | |
| free_slots_.reserve(n_small_); | |
| for (size_t i = n_small_; i-- > 0;) { | |
| free_slots_.push_back((int32_t) i); | |
| } | |
| free_big_slots_.reserve(n_big_); | |
| for (size_t i = n_slots_; i-- > n_small_;) { | |
| free_big_slots_.push_back((int32_t) i); | |
| } | |
| used_.store(0); | |
| hist_.assign((size_t) info_.n_layer, {}); | |
| for (auto & per_layer : hist_) { | |
| per_layer.assign((size_t) info_.n_expert, {}); | |
| } | |
| prev_used_.assign((size_t) info_.n_layer, {}); | |
| cur_used_.assign((size_t) info_.n_layer, {}); | |
| } | |
| void pager::init(const config & cfg, ggml_context * ctx) { | |
| (void) ctx; | |
| if (ready_) { | |
| return; | |
| } | |
| cfg_ = cfg; | |
| if (cfg_.model_path.empty()) { | |
| const char * v = getenv("SDQ_MODEL_PATH"); | |
| cfg_.model_path = v ? v : ""; | |
| } | |
| if (cfg_.model_path.empty()) { | |
| throw std::runtime_error("SDQ 缺少模型路徑(設定 SDQ_MODEL_PATH)"); | |
| } | |
| open_model_fd(); | |
| if (fd_ < 0) { | |
| throw std::runtime_error("無法開啟模型檔:" + cfg_.model_path); | |
| } | |
| load_gguf_header(cfg_.model_path); | |
| // 順序有意義:I/O 池的對齊緩衝必須在 alloc_arena() 之前配置, | |
| // 否則 arena 會把它當成「不存在的 RAM」,實際上就偷用了 8 GB 預算。 | |
| io_pool_start(); | |
| alloc_arena(); | |
| if (cfg_.prefetch) { | |
| pf_ = std::thread([this] { pf_worker(); }); | |
| } | |
| ready_ = true; | |
| print_summary(); | |
| } | |
| const layer_desc * pager::layer(int il) const { | |
| if (il < 0 || il >= (int) layers_.size()) { | |
| return nullptr; | |
| } | |
| const layer_desc & L = layers_[(size_t) il]; | |
| if (L.experts.empty()) { | |
| return nullptr; | |
| } | |
| return &L; | |
| } | |
| void pager::open_model_fd() { | |
| fd_ = open(cfg_.model_path.c_str(), O_RDONLY); | |
| if (fd_ < 0) { | |
| return; | |
| } | |
| // O_DIRECT:讀資料完全不進 kernel page cache。 | |
| // 這是「RAM 帳目誠實」的關鍵 —— 否則核心預讀會讓 mmap 檔案的頁算進 RSS | |
| // (實測 4K prefill 時檔案對映 RSS 可以多到 6 GB 以上),而且量到的 | |
| // 「SSD 讀取量」會被上一輪的 page cache 灌水。 | |
| const char * no_direct = getenv("SDQ_NO_DIRECT_IO"); | |
| if (no_direct && *no_direct && *no_direct != '0') { | |
| return; | |
| } | |
| fd_direct_ = open(cfg_.model_path.c_str(), O_RDONLY | O_DIRECT); | |
| if (fd_direct_ >= 0) { | |
| use_direct_ = true; | |
| fprintf(stderr, "[sdq] 使用 O_DIRECT 讀取模型檔(不佔用 kernel page cache)\n"); | |
| } else { | |
| fprintf(stderr, "[sdq] O_DIRECT 不可用,改用 buffered + fadvise(DONTNEED)\n"); | |
| } | |
| } | |
| // 有對齊的緩衝區(O_DIRECT 要求),thread_local 避免競爭 | |
| static thread_local uint8_t * g_direct_buf = nullptr; | |
| static thread_local size_t g_direct_sz = 0; | |
| // 與 read_range_direct 相同,但用呼叫端提供的緩衝區(已配置好、已算進 RSS)。 | |
| // I/O 執行緒池走這條:每條 worker 在啟動時就配置好對齊緩衝,於是 arena sizing 時 | |
| // 這筆記憶體已經在 RSS 裡,不會變成「額外偷用的 RAM」。 | |
| bool pager::read_range_staged(uint64_t off, uint8_t * dst, size_t size, | |
| uint8_t * stage, size_t stage_sz) { | |
| if (!use_direct_) { | |
| return read_range(off, dst, size); | |
| } | |
| const size_t blk = 4096; | |
| const uint64_t start = off & ~(uint64_t)(blk - 1); | |
| const uint64_t end = (off + size + blk - 1) & ~(uint64_t)(blk - 1); | |
| const size_t bytes = (size_t)(end - start); | |
| if (stage_sz < bytes) { | |
| // 預先配置不足(不該發生):退回 thread_local 的通用路徑 | |
| return read_range_direct(off, dst, size); | |
| } | |
| size_t done = 0; | |
| while (done < bytes) { | |
| const ssize_t n = pread(fd_direct_, stage + done, bytes - done, | |
| (off_t)(start + done)); | |
| if (n <= 0) { | |
| return false; | |
| } | |
| done += (size_t) n; | |
| } | |
| memcpy(dst, stage + (off - start), size); | |
| return true; | |
| } | |
| bool pager::read_range_direct(uint64_t off, uint8_t * dst, size_t size) { | |
| if (!use_direct_) { | |
| return read_range(off, dst, size); | |
| } | |
| const size_t blk = 4096; | |
| const uint64_t start = off & ~(uint64_t)(blk - 1); | |
| const uint64_t end = (off + size + blk - 1) & ~(uint64_t)(blk - 1); | |
| const size_t bytes = (size_t)(end - start); | |
| if (g_direct_sz < bytes) { | |
| if (g_direct_buf) { | |
| free(g_direct_buf); | |
| } | |
| if (posix_memalign((void **) &g_direct_buf, blk, bytes) != 0) { | |
| g_direct_buf = nullptr; | |
| g_direct_sz = 0; | |
| use_direct_ = false; | |
| return read_range(off, dst, size); | |
| } | |
| g_direct_sz = bytes; | |
| } | |
| size_t done = 0; | |
| while (done < bytes) { | |
| const ssize_t n = pread(fd_direct_, g_direct_buf + done, bytes - done, | |
| (off_t)(start + done)); | |
| if (n <= 0) { | |
| return false; | |
| } | |
| done += (size_t) n; | |
| } | |
| memcpy(dst, g_direct_buf + (off - start), size); | |
| return true; | |
| } | |
| bool pager::read_range(uint64_t off, uint8_t * dst, size_t size) { | |
| size_t done = 0; | |
| while (done < size) { | |
| const ssize_t n = pread(fd_, dst + done, size - done, (off_t) (off + done)); | |
| if (n <= 0) { | |
| return false; | |
| } | |
| done += (size_t) n; | |
| } | |
| // 讓 RAM 帳目誠實:把 kernel 對這個檔案的 page cache 丟掉,避免「藏在核心裡的 RAM」 | |
| if (cfg_.fadvise_dontneed) { | |
| posix_fadvise(fd_, (off_t) off, (off_t) size, POSIX_FADV_DONTNEED); | |
| } | |
| return true; | |
| } | |
| bool pager::read_expert(int il, int ie, uint8_t * dst) { | |
| const layer_desc * L = layer(il); | |
| if (!L || ie < 0 || ie >= (int) L->experts.size()) { | |
| return false; | |
| } | |
| const expert_desc & ed = L->experts[(size_t) ie]; | |
| for (int p = 0; p < 3; p++) { | |
| if (ed.seg[p].size == 0) { | |
| continue; | |
| } | |
| if (!read_range_direct(ed.seg[p].off, dst + ed.slot_off[p], (size_t) ed.seg[p].size)) { | |
| return false; | |
| } | |
| } | |
| return true; | |
| } | |
| // ---------------------------------------------------------------- I/O 執行緒池 | |
| // | |
| // 問題:一個 expert = gate / up / down 三段,檔案上是三段不相鄰的資料。 | |
| // 循序讀就是 3 次 pread 排在同一條執行緒上,而單次 1.9 MB 的 pread 在本機量到 | |
| // ~7 ms(≈275 MB/s 的單流頻寬)。因此「一個 expert 要多久」= 3 × 單次延遲, | |
| // 這條路徑的長度跟「有多少條執行緒」無關 —— mode C 再怎麼排程都救不了。 | |
| // | |
| // 做法:把三段交給三條執行緒同時 pread,臨界路徑降到約「單段」的時間。 | |
| // 量測(8 個 expert 全部 miss 的一層,各 9 次取中位數): | |
| // 8 執行緒 × 三段序列讀 : 12.60 ms | |
| // 16 執行緒 × 三段平行讀 : 7.57 ms | |
| // 24 執行緒 × 三段平行讀 : 7.58 ms ← 16 就飽和,所以預設 3 × n_threads | |
| // | |
| // 這個池子對 prefill 更有用:一層要 256 個 expert = 768 個讀取工作, | |
| // 序列讀時 16 執行緒只能拿到 16 路的頻寬 utilization。 | |
| void pager::io_pool_start() { | |
| if (!cfg_.io_parallel) { | |
| fprintf(stderr, "[sdq] I/O 段平行:關(每個 expert 的 gate/up/down 循序讀)\n"); | |
| return; | |
| } | |
| int n = cfg_.io_threads > 0 ? cfg_.io_threads : cfg_.n_threads * 3; | |
| if (n < 4) { | |
| n = 4; | |
| } | |
| if (n > 48) { | |
| n = 48; | |
| } | |
| // 先算「一段權重最大有多少位元組」,每條 worker 配一塊對齊緩衝。 | |
| // 這必须在 alloc_arena() 之前做完,否則這筆 RAM 會變成帳外支出。 | |
| size_t max_seg = 0; | |
| for (const auto & L : layers_) { | |
| for (const auto & ed : L.experts) { | |
| for (int p = 0; p < 3; p++) { | |
| max_seg = std::max(max_seg, (size_t) ed.seg[p].size); | |
| } | |
| } | |
| } | |
| const size_t stage_bytes = ((max_seg + 2 * 4096) + 4095) & ~(size_t) 4095; | |
| io_stage_.assign((size_t) n, io_stage{}); | |
| for (int i = 0; i < n; i++) { | |
| void * p = nullptr; | |
| if (posix_memalign(&p, 4096, stage_bytes) == 0 && p) { | |
| memset(p, 0, stage_bytes); // 實際碰頁,讓它出現在 RSS 裡 | |
| io_stage_[(size_t) i].buf = (uint8_t *) p; | |
| io_stage_[(size_t) i].sz = stage_bytes; | |
| } | |
| } | |
| io_split_ = cfg_.io_split > 0 ? cfg_.io_split : 1; | |
| io_stop_ = false; | |
| io_pool_.reserve((size_t) n); | |
| for (int i = 0; i < n; i++) { | |
| const size_t wid = (size_t) i; | |
| try { | |
| io_pool_.emplace_back([this, wid] { io_worker(wid); }); | |
| } catch (...) { | |
| break; | |
| } | |
| } | |
| io_threads_ = (int) io_pool_.size(); | |
| io_pool_on_ = io_threads_ > 0; | |
| if (io_pool_on_) { | |
| fprintf(stderr, "[sdq] I/O 段平行:%d 條讀取執行緒 × %zu KB 對齊緩衝" | |
| "(gate/up/down 同時讀,每段再切 %d 塊)\n", | |
| io_threads_, stage_bytes >> 10, io_split_); | |
| } | |
| } | |
| void pager::io_pool_stop() { | |
| if (io_pool_.empty()) { | |
| for (auto & s : io_stage_) { | |
| free(s.buf); | |
| s.buf = nullptr; | |
| } | |
| io_stage_.clear(); | |
| return; | |
| } | |
| { | |
| std::lock_guard<std::mutex> lk(io_mu_); | |
| io_stop_ = true; | |
| } | |
| io_cv_.notify_all(); | |
| for (auto & t : io_pool_) { | |
| if (t.joinable()) { | |
| t.join(); | |
| } | |
| } | |
| io_pool_.clear(); | |
| io_pool_on_ = false; | |
| for (auto & s : io_stage_) { | |
| free(s.buf); | |
| s.buf = nullptr; | |
| s.sz = 0; | |
| } | |
| io_stage_.clear(); | |
| } | |
| void pager::io_worker(size_t wid) { | |
| for (;;) { | |
| io_job job{}; | |
| { | |
| std::unique_lock<std::mutex> lk(io_mu_); | |
| io_cv_.wait(lk, [this] { return io_stop_ || !io_q_.empty(); }); | |
| if (io_q_.empty()) { | |
| if (io_stop_) { | |
| return; | |
| } | |
| continue; | |
| } | |
| job = io_q_.front(); | |
| io_q_.pop_front(); | |
| } | |
| io_jobs_.fetch_add(1, std::memory_order_relaxed); | |
| // 讀取本身不持鎖:io_mu_ 只保護佇列與完成計數 | |
| const bool ok = (wid < io_stage_.size() && io_stage_[wid].buf) | |
| ? read_range_staged(job.off, job.dst, job.size, | |
| io_stage_[wid].buf, io_stage_[wid].sz) | |
| : read_range_direct(job.off, job.dst, job.size); | |
| { | |
| std::lock_guard<std::mutex> lk(io_mu_); | |
| if (!ok) { | |
| job.ctx->ok = false; | |
| } | |
| if (job.ctx->remaining.fetch_sub(1) == 1) { | |
| io_done_cv_.notify_all(); | |
| } | |
| } | |
| } | |
| } | |
| // 把 (il,ie) 的三段交給執行緒池並行讀,呼叫端等到整批完成才返回。 | |
| // ctx 放在呼叫端 stack:worker 只會在 remaining 歸零之後才再碰它, | |
| // 而我們正是等到 remaining==0 才離開函式,所以沒有 use-after-return。 | |
| bool pager::read_expert_par(int il, int ie, uint8_t * dst) { | |
| const layer_desc * L = layer(il); | |
| if (!L || ie < 0 || ie >= (int) L->experts.size()) { | |
| return false; | |
| } | |
| if (!io_pool_on_) { | |
| return read_expert(il, ie, dst); | |
| } | |
| const expert_desc & ed = L->experts[(size_t) ie]; | |
| io_ctx ctx; | |
| int n = 0; | |
| // 把每一段再切成一樣大的好幾塊,交給不同的執行緒同時讀。 | |
| // | |
| // 為什麼:單一 pread 的時間 ≈ 固定延遲 + 位元組數 / 單流頻寬。 | |
| // 本機實測單流約 275 MB/s、4K 隨機讀延遲約 0.57 ms,所以 | |
| // 590 KB 讀一次 ≈ 0.57 + 2.1 = 2.7 ms | |
| // 590 KB 切成 2 塊 ≈ 0.57 + 1.05 = 1.6 ms | |
| // 而「讀完一個 expert」的時間正是每層的關鍵路徑,切塊就是直接砍關鍵路徑。 | |
| // 設定檔 SDQ_IO_SPLIT。**實測預設是 1(不切塊)**:切塊在 24 條讀取執行緒 | |
| // 之下沒有可重現的改善(split=1..6 的中位數在 4.9~5.9 之間跳動,純噪音), | |
| // 因為 8 個 expert × 3 段 = 24 個並行請求已經把裝置的佇列灌飽, | |
| // 這時瓶頸是裝置總頻寬而不是單一請求的延遲。保留這個開關是為了其他裝置。 | |
| const int split = io_split_; | |
| const size_t blk = 4096; | |
| { | |
| std::lock_guard<std::mutex> lk(io_mu_); | |
| for (int p = 0; p < 3; p++) { | |
| const size_t sz = (size_t) ed.seg[p].size; | |
| if (sz == 0) { | |
| continue; | |
| } | |
| uint8_t * base = dst + ed.slot_off[p]; | |
| if (split <= 1 || sz <= 2 * blk) { | |
| io_q_.push_back(io_job{ed.seg[p].off, base, sz, &ctx}); | |
| n++; | |
| continue; | |
| } | |
| size_t chunk = ((sz + (size_t) split - 1) / (size_t) split + blk - 1) & ~(blk - 1); | |
| for (size_t off = 0; off < sz; off += chunk) { | |
| const size_t take = sz - off < chunk ? sz - off : chunk; | |
| io_q_.push_back(io_job{ed.seg[p].off + off, base + off, take, &ctx}); | |
| n++; | |
| } | |
| } | |
| ctx.remaining.store(n); | |
| } | |
| if (n == 0) { | |
| return true; | |
| } | |
| io_cv_.notify_all(); | |
| const uint64_t t0 = now_us(); | |
| { | |
| std::unique_lock<std::mutex> lk(io_mu_); | |
| io_done_cv_.wait(lk, [&ctx] { return ctx.remaining.load() == 0; }); | |
| } | |
| io_wait_ns_.fetch_add((now_us() - t0) * 1000ull, std::memory_order_relaxed); | |
| return ctx.ok; | |
| } | |
| uint64_t pager::read_expert_span(int il, const int32_t * experts, int n, const int32_t * slots) { | |
| const layer_desc * L = layer(il); | |
| if (!L || n <= 0) { | |
| return 0; | |
| } | |
| const uint64_t t_start = now_us(); | |
| // 每個執行緒一份暫存區(多執行緒會同時呼叫,不能共用一個 buffer) | |
| static thread_local std::vector<uint8_t> scratch_buf; | |
| uint64_t total = 0; | |
| // 把 experts 切成「檔案編號連續」的區段 | |
| int i = 0; | |
| while (i < n) { | |
| int j = i + 1; | |
| while (j < n && experts[j] == experts[j - 1] + 1) { | |
| j++; | |
| } | |
| const int cnt = j - i; | |
| const expert_desc & first = L->experts[(size_t) experts[i]]; | |
| const expert_desc & last = L->experts[(size_t) experts[j - 1]]; | |
| for (int p = 0; p < 3; p++) { | |
| if (first.seg[p].size == 0 || last.seg[p].size == 0) { | |
| continue; | |
| } | |
| const uint64_t off = first.seg[p].off; | |
| const uint64_t end = last.seg[p].off + last.seg[p].size; | |
| const size_t bytes = (size_t) (end - off); | |
| if (scratch_buf.size() < bytes) { | |
| scratch_buf.assign(bytes, 0); | |
| } | |
| if (!read_range_direct(off, scratch_buf.data(), bytes)) { | |
| return total; | |
| } | |
| total += bytes; | |
| // 散到各自的槽位:scratch 從 off 開始,所以偏移是「相對於 off」 | |
| for (int k = i; k < j; k++) { | |
| const expert_desc & ed = L->experts[(size_t) experts[k]]; | |
| const size_t seg_bytes = (size_t) ed.seg[p].size; | |
| const size_t delta = (size_t) (ed.seg[p].off - off); | |
| if (slots[k] < 0) { | |
| continue; | |
| } | |
| if ((size_t) slots[k] >= n_slots_ || (uint64_t) ed.slot_off[p] + seg_bytes > slot_bytes_ | |
| || delta + seg_bytes > scratch_buf.size()) { | |
| fprintf(stderr, "[sdq-pager] span 越界:il=%d slot=%d/%zu slot_off=%lld seg=%zu delta=%zu scratch=%zu p=%d e=%d\n", | |
| il, slots[k], n_slots_, (long long) ed.slot_off[p], seg_bytes, delta, scratch_buf.size(), p, experts[k]); | |
| return total; | |
| } | |
| memcpy(slot_ptr(slots[k]) + ed.slot_off[p], | |
| scratch_buf.data() + delta, seg_bytes); | |
| } | |
| } | |
| i = j; | |
| } | |
| stats_.demand_wait_us.fetch_add(now_us() - t_start, std::memory_order_relaxed); | |
| return total; | |
| } | |
| int pager::reserve_slot(int il, int ie) { | |
| if (!ready_) { | |
| return -1; | |
| } | |
| count_request(il, ie); | |
| std::lock_guard<std::mutex> lk(mu_); | |
| const int32_t exist = find_slot_locked(il, ie); | |
| if (exist >= 0) { | |
| slots_[(size_t) exist].last_use = ++clock_; | |
| slots_[(size_t) exist].pin++; | |
| touch_freq_locked(slots_[(size_t) exist]); | |
| return exist; | |
| } | |
| const int32_t s = victim_locked(); | |
| if (s < 0) { | |
| return -1; | |
| } | |
| slot_state & st = slots_[(size_t) s]; | |
| if (st.layer >= 0) { | |
| page_.erase(((int64_t) st.layer << 20) | (int64_t) st.expert); | |
| stats_.evictions.fetch_add(1, std::memory_order_relaxed); | |
| used_.fetch_sub(1, std::memory_order_relaxed); | |
| } | |
| st.layer = il; | |
| st.expert = ie; | |
| st.last_use = ++clock_; | |
| st.freq = 1; | |
| st.pin = 1; | |
| st.is_hot = hot_lookup(il, ie); | |
| page_[((int64_t) il << 20) | (int64_t) ie] = s; | |
| used_.fetch_add(1, std::memory_order_relaxed); | |
| stats_.expert_misses.fetch_add(1, std::memory_order_relaxed); | |
| return s; | |
| } | |
| int pager::probe(int il, int ie) { | |
| if (!ready_) { | |
| return -1; | |
| } | |
| stats_.expert_requests.fetch_add(1, std::memory_order_relaxed); | |
| count_request(il, ie); | |
| std::lock_guard<std::mutex> lk(mu_); | |
| const int32_t s = find_slot_locked(il, ie); | |
| if (s < 0) { | |
| stats_.expert_misses.fetch_add(1, std::memory_order_relaxed); | |
| return -1; | |
| } | |
| slots_[(size_t) s].last_use = ++clock_; | |
| slots_[(size_t) s].pin++; | |
| touch_freq_locked(slots_[(size_t) s]); | |
| stats_.expert_hits.fetch_add(1, std::memory_order_relaxed); | |
| return s; | |
| } | |
| void pager::note_span_read(uint64_t bytes) { | |
| stats_.ssb_read_bytes_demand.fetch_add(bytes, std::memory_order_relaxed); | |
| stats_.demand_wait_us.fetch_add(0, std::memory_order_relaxed); | |
| } | |
| // 衰減式頻率:每 8192 次使用就把所有計數減半,讓「過去熱門、現在變冷」的 expert 能被淘汰 | |
| void pager::touch_freq_locked(slot_state & st) { | |
| if (++freq_tick_ >= 8192) { | |
| freq_tick_ = 0; | |
| for (auto & s : slots_) { | |
| s.freq /= 2; | |
| } | |
| } | |
| if (st.freq < 0xFFFFFFFFu) { | |
| st.freq++; | |
| } | |
| } | |
| // ============================== OPT-1:熱門 expert 剖面 ============================== | |
| // 問題實測(docs/03):arena 3.88 GiB 只放得下 2044 個 expert,總共卻有 10240 個, | |
| // prefill 的幾百個「只用一次」的專家會把熱區整個洗掉,命中率卡在 ~70%。 | |
| // 解法借用 Strata 的做法(src/core/expert_cache.cpp):**離線剖面 + 靜態排行榜 + | |
| // 熱頁永不淘汰**。這裡不動檔案格式,用純文字存排行榜。 | |
| void pager::count_request(int il, int ie) { | |
| if (prof_.empty() || il < 0 || ie < 0) { | |
| return; | |
| } | |
| const size_t idx = (size_t) il * (size_t) info_.n_expert + (size_t) ie; | |
| if (idx < prof_.size()) { | |
| // 單執行緒 ++ 即可:請求次數是統計用途,不參與任何正確性判斷 | |
| if (prof_[idx] < 0xFFFFFFFFu) { | |
| prof_[idx]++; | |
| } | |
| } | |
| } | |
| bool pager::hot_lookup(int il, int ie) const { | |
| if (hot_.empty() || il < 0 || ie < 0) { | |
| return false; | |
| } | |
| const size_t idx = (size_t) il * (size_t) info_.n_expert + (size_t) ie; | |
| return idx < hot_.size() && hot_[idx] != 0; | |
| } | |
| // 挑一個容得下 (il,ie) 的槽。與 victim_locked 的差別:**先只在冷頁裡找**, | |
| // 找不到才退讓去動熱頁。這就是「熱區不會被 prefill 洗掉」的機制。 | |
| // 這個槽位放得下「需要 big 的 expert」嗎? | |
| // 小槽位:只能放小 expert | |
| // 大槽位:小expert 也放得下(只是浪費一點空間) | |
| static inline bool slot_accepts(bool need_big, bool slot_is_big) { | |
| return !need_big || slot_is_big; | |
| } | |
| int pager::admit_slot_locked(int il, int ie) { | |
| // il/ie 都 < 0 代表「不指定大小」→ 兩種槽位都可以(舊的呼叫路徑)。 | |
| const bool need_big = (il >= 0 && ie >= 0) ? need_big_slot(il, ie) : false; | |
| // 完全沒用過的空槽永遠優先(它不是熱頁,佔用等於浪費) | |
| // 大 expert 先看大槽位;小 expert 先看小槽位。 | |
| std::vector<int32_t> & first = need_big ? free_big_slots_ : free_slots_; | |
| std::vector<int32_t> & second = need_big ? free_slots_ : free_big_slots_; | |
| while (!first.empty()) { | |
| const int32_t s = first.back(); | |
| if (!slots_[(size_t) s].busy) { | |
| first.pop_back(); | |
| return s; | |
| } | |
| } | |
| // 第一輪:只在「冷頁」裡挑最久未用的。 | |
| // want_hot=true → 熱門 expert 進來時,把冷頁換掉(不動其他熱頁) | |
| // want_hot=false → 冷門 expert 進來時,優先吃冷頁 | |
| // 兩種情況的規則相同,所以 hot_lookup 的結果不影響選擇順序,只影響標記。 | |
| for (int pass = 0; pass < 2; pass++) { | |
| int64_t best = -1; | |
| uint64_t best_clock = UINT64_MAX; | |
| for (size_t s = 0; s < n_slots_; s++) { | |
| const slot_state & st = slots_[s]; | |
| if (st.pin > 0 || st.busy) { | |
| continue; | |
| } | |
| if (pass == 0 && st.is_hot) { | |
| continue; // 第一輪不動熱頁 | |
| } | |
| if (!slot_accepts(need_big, is_big_slot(s))) { | |
| continue; | |
| } | |
| if (st.last_use < best_clock) { | |
| best_clock = st.last_use; | |
| best = (int64_t) s; | |
| } | |
| } | |
| if (best >= 0) { | |
| return (int) best; | |
| } | |
| } | |
| // 退讓:另一種大小的空槽。**一定要檢查放不放得下** —— 小槽位放不下大 expert, | |
| // 硬塞會讓 read_expert 寫出槽位邊界(靜默把別的 expert 寫壞)。 | |
| while (!second.empty()) { | |
| const int32_t s = second.back(); | |
| if (!slots_[(size_t) s].busy && | |
| slot_accepts(need_big, is_big_slot(s))) { | |
| second.pop_back(); | |
| return s; | |
| } | |
| if (!slots_[(size_t) s].busy) { | |
| second.pop_back(); | |
| } else { | |
| break; | |
| } | |
| } | |
| return -1; | |
| } | |
| bool pager::load_profile(const std::string & path) { | |
| std::FILE * f = std::fopen(path.c_str(), "rb"); | |
| if (!f) { | |
| fprintf(stderr, "[sdq] 找不到 expert 剖面 %s(將以動態 LRU/LFU 運作)\n", path.c_str()); | |
| return false; | |
| } | |
| char magic[8] = {0}; | |
| if (std::fread(magic, 1, 4, f) != 4 || std::memcmp(magic, "SDQP", 4) != 0) { | |
| std::fclose(f); | |
| fprintf(stderr, "[sdq] %s 不是 SDQ expert 剖面(缺少 SDQP 標頭)\n", path.c_str()); | |
| return false; | |
| } | |
| uint32_t hdr[3] = {0, 0, 0}; | |
| if (std::fread(hdr, sizeof(uint32_t), 3, f) != 3) { | |
| std::fclose(f); | |
| return false; | |
| } | |
| const uint32_t n = hdr[2]; | |
| std::vector<std::pair<uint32_t, std::pair<int32_t, int32_t>>> ranked; | |
| ranked.reserve(n); | |
| for (uint32_t i = 0; i < n; i++) { | |
| int32_t l = 0, e = 0; | |
| uint32_t c = 0; | |
| if (std::fread(&l, sizeof(int32_t), 1, f) != 1) break; | |
| if (std::fread(&e, sizeof(int32_t), 1, f) != 1) break; | |
| if (std::fread(&c, sizeof(uint32_t), 1, f) != 1) break; | |
| ranked.push_back({c, {l, e}}); | |
| } | |
| std::fclose(f); | |
| std::sort(ranked.begin(), ranked.end(), | |
| [](const auto & a, const auto & b) { return a.first > b.first; }); | |
| hot_.assign((size_t) info_.n_layer * (size_t) info_.n_expert, 0); | |
| const int want = cfg_.hot_pin > 0 ? cfg_.hot_pin : (int) ranked.size(); | |
| n_hot_ = 0; | |
| for (const auto & r : ranked) { | |
| if (n_hot_ >= want) { | |
| break; | |
| } | |
| const int l = r.second.first, e = r.second.second; | |
| if (l < 0 || l >= info_.n_layer || e < 0 || e >= info_.n_expert) { | |
| continue; | |
| } | |
| const size_t idx = (size_t) l * (size_t) info_.n_expert + (size_t) e; | |
| if (!hot_[idx]) { | |
| hot_[idx] = 1; | |
| n_hot_++; | |
| } | |
| } | |
| fprintf(stderr, "[sdq] expert 剖面:%s → 釘選 %d / %zu 個熱門 expert(永不淘汰)\n", | |
| path.c_str(), n_hot_, ranked.size()); | |
| return n_hot_ > 0; | |
| } | |
| void pager::dump_profile(const std::string & path) { | |
| if (prof_.empty()) { | |
| return; | |
| } | |
| std::vector<std::pair<uint32_t, std::pair<int32_t, int32_t>>> ranked; | |
| ranked.reserve(prof_.size()); | |
| uint64_t total = 0; | |
| for (int l = 0; l < info_.n_layer; l++) { | |
| for (int e = 0; e < info_.n_expert; e++) { | |
| const uint32_t c = prof_[(size_t) l * (size_t) info_.n_expert + (size_t) e]; | |
| if (c > 0) { | |
| ranked.push_back({c, {l, e}}); | |
| total += c; | |
| } | |
| } | |
| } | |
| std::sort(ranked.begin(), ranked.end(), | |
| [](const auto & a, const auto & b) { return a.first > b.first; }); | |
| std::FILE * f = std::fopen(path.c_str(), "wb"); | |
| if (!f) { | |
| fprintf(stderr, "[sdq] 無法寫出 expert 剖面 %s\n", path.c_str()); | |
| return; | |
| } | |
| std::fwrite("SDQP", 1, 4, f); | |
| const uint32_t hdr[3] = {1, (uint32_t) info_.n_layer, (uint32_t) ranked.size()}; | |
| std::fwrite(hdr, sizeof(uint32_t), 3, f); | |
| for (const auto & r : ranked) { | |
| const int32_t l = r.second.first, e = r.second.second; | |
| const uint32_t c = r.first; | |
| std::fwrite(&l, sizeof(int32_t), 1, f); | |
| std::fwrite(&e, sizeof(int32_t), 1, f); | |
| std::fwrite(&c, sizeof(uint32_t), 1, f); | |
| } | |
| std::fclose(f); | |
| // 覆盖率報告:前 K 名可以蓋掉多少請求 → 直接決定 hot_pin 該設多少 | |
| uint64_t acc = 0; | |
| const size_t ks[6] = {512, 1024, 1536, 2048, 3072, 4096}; | |
| fprintf(stderr, "[sdq] expert 剖面已寫到 %s(%zu 個有被用到,總請求 %llu)\n", | |
| path.c_str(), ranked.size(), (unsigned long long) total); | |
| size_t ki = 0; | |
| for (size_t i = 0; i < ranked.size(); i++) { | |
| acc += ranked[i].first; | |
| if (ki < 6 && i + 1 == ks[ki]) { | |
| fprintf(stderr, "[sdq] 前 %5zu 名涵蓋 %6.2f%% 的請求\n", | |
| ks[ki], 100.0 * (double) acc / (double) (total ? total : 1)); | |
| ki++; | |
| } | |
| } | |
| } | |
| int pager::find_slot_locked(int il, int ie) { | |
| const int64_t key = ((int64_t) il << 20) | (int64_t) ie; | |
| auto it = page_.find(key); | |
| return it == page_.end() ? -1 : it->second; | |
| } | |
| // 找一個「可以放 (il,ie) 而且 last_use <= stale_before」的槽位。 | |
| // | |
| // 為什麼不能直接用 victim_stale_locked():大小兩種槽位之後,那個函式會掃過**所有** | |
| // 槽位(包含小槽位),於是一個需要大槽位的 expert 可能拿到小槽位, | |
| // 預取時 read_expert 會寫出槽位邊界 → 靜默把鄰居槽位的資料寫壞。 | |
| // 這是引進大小兩種槽位時真正會咬人的地方。 | |
| // | |
| // 順序:先看對應尺寸的空槽(從沒用過,最該拿來放預取), | |
| // 再掃全部槽位挑最久未用的(放得下的那種尺寸)。 | |
| int pager::victim_stale_for(int il, int ie, uint64_t stale_before) { | |
| const bool need_big = need_big_slot(il, ie); | |
| std::vector<int32_t> & first = need_big ? free_big_slots_ : free_slots_; | |
| std::vector<int32_t> & second = need_big ? free_slots_ : free_big_slots_; | |
| for (std::vector<int32_t> * lst : { &first, &second }) { | |
| while (!lst->empty()) { | |
| const int32_t s = lst->back(); | |
| if (slots_[(size_t) s].busy) { | |
| break; | |
| } | |
| if (slots_[(size_t) s].last_use <= stale_before && | |
| slot_accepts(need_big, is_big_slot(s))) { | |
| lst->pop_back(); | |
| return s; | |
| } | |
| lst->pop_back(); | |
| } | |
| } | |
| int64_t best = -1; | |
| uint64_t best_clock = UINT64_MAX; | |
| for (size_t i = 0; i < n_slots_; i++) { | |
| const slot_state & st = slots_[i]; | |
| if (st.pin > 0 || st.busy || st.last_use > stale_before) { | |
| continue; | |
| } | |
| if (!slot_accepts(need_big, is_big_slot(i))) { | |
| continue; | |
| } | |
| if (st.last_use < best_clock) { | |
| best_clock = st.last_use; | |
| best = (int64_t) i; | |
| } | |
| } | |
| return (int) best; | |
| } | |
| int pager::victim_stale_locked(uint64_t stale_before) { | |
| // 只挑「last_use 很舊」而且沒被 pin 的槽:保護最近的工作集 | |
| while (!free_slots_.empty()) { | |
| const int32_t s = free_slots_.back(); | |
| if (!slots_[(size_t) s].busy && slots_[(size_t) s].last_use <= stale_before) { | |
| free_slots_.pop_back(); | |
| return s; | |
| } | |
| if (free_slots_.size() * 4 < n_slots_) { | |
| break; | |
| } | |
| free_slots_.pop_back(); | |
| } | |
| int64_t best = -1; | |
| uint64_t best_clock = UINT64_MAX; | |
| for (size_t i = 0; i < n_slots_; i++) { | |
| const slot_state & st = slots_[i]; | |
| if (st.pin > 0 || st.busy || st.last_use > stale_before) { | |
| continue; | |
| } | |
| if (st.last_use < best_clock) { | |
| best_clock = st.last_use; | |
| best = (int64_t) i; | |
| } | |
| } | |
| return (int) best; | |
| } | |
| int pager::acquire(int il, int ie, bool blocking) { | |
| if (!ready_) { | |
| return -1; | |
| } | |
| std::unique_lock<std::mutex> lk(mu_); | |
| int32_t s = find_slot_locked(il, ie); | |
| if (s >= 0 && slots_[s].busy) { | |
| // 這個槽位**正在被填入**:另一條執行緒剛 admit 它、資料還沒讀完。 | |
| // | |
| // 這裡一定要等,不能當成命中直接用。否則讀到的是「寫到一半的權重」, | |
| // 而且那個錯誤是**非決定性**的 —— 每次執行結果可能不同(實測在 | |
| // prefill 時會讓同一個提示詞跑出不同的文字)。 | |
| // busy 的槽位在 admit/victim 那邊本來就會被跳過,所以這裡是唯一 | |
| // 還能拿到它的路徑,也因此必須在這裡等。 | |
| const uint64_t t0 = now_us(); | |
| while (slots_[s].busy) { | |
| lk.unlock(); | |
| std::this_thread::sleep_for(std::chrono::microseconds(200)); | |
| lk.lock(); | |
| if (now_us() - t0 > (uint64_t) env_size("SDQ_BUSY_WAIT_MS", 30000) * 1000) { | |
| fprintf(stderr, | |
| "[sdq] 致命:等待槽位 il=%d ie=%d 讀入超過 %u ms(busy 一直沒解除)\n", | |
| il, ie, (unsigned) env_size("SDQ_BUSY_WAIT_MS", 30000)); | |
| std::fflush(stderr); | |
| abort(); | |
| } | |
| } | |
| // 讀取可能失敗(那時槽位已被清空)→ 重新走一次正常的 miss 路徑 | |
| if (find_slot_locked(il, ie) != s) { | |
| s = -1; | |
| } | |
| } | |
| if (s >= 0) { | |
| slots_[s].last_use = ++clock_; | |
| slots_[s].pin++; | |
| touch_freq_locked(slots_[(size_t) s]); | |
| stats_.expert_requests.fetch_add(1, std::memory_order_relaxed); | |
| stats_.expert_hits.fetch_add(1, std::memory_order_relaxed); | |
| return s; | |
| } | |
| stats_.expert_requests.fetch_add(1, std::memory_order_relaxed); | |
| count_request(il, ie); | |
| stats_.expert_misses.fetch_add(1, std::memory_order_relaxed); | |
| // 先找一個已經在背景讀入這頁的槽 | |
| for (size_t i = 0; i < n_slots_; i++) { | |
| slot_state & st = slots_[i]; | |
| if (st.busy && st.layer == il && st.expert == ie) { | |
| // 等背景執行緒讀完 | |
| lk.unlock(); | |
| const uint64_t t0 = now_us(); | |
| while (true) { | |
| lk.lock(); | |
| if (!slots_[i].busy) { | |
| break; | |
| } | |
| lk.unlock(); | |
| std::this_thread::sleep_for(std::chrono::microseconds(200)); | |
| } | |
| s = (int32_t) i; | |
| lk.unlock(); | |
| stats_.prefetch_wait_us.fetch_add(now_us() - t0, std::memory_order_relaxed); | |
| std::unique_lock<std::mutex> lk2(mu_); | |
| slots_[s].last_use = ++clock_; | |
| slots_[s].pin++; | |
| stats_.expert_hits.fetch_add(1, std::memory_order_relaxed); | |
| stats_.prefetch_used.fetch_add(1, std::memory_order_relaxed); | |
| return s; | |
| } | |
| } | |
| if (!blocking) { | |
| return -1; | |
| } | |
| // 選槽位時必須知道「這個 expert 需要哪一種大小」,否則大 expert 可能被塞進小槽位 | |
| // (寫出邊界 → 靜默算錯),或者小 expert 佔掉大槽位(浪費 6.8% 的 arena)。 | |
| s = admit_slot_locked(il, ie); | |
| if (s < 0) { | |
| // 拿不到槽位:短暫等待別的執行緒解除 pin(每層開頭會 unpin_all)。 | |
| // 這裡「重試」而不是直接回 -1,因為回 -1 會讓呼叫端**跳過這個 expert**, | |
| // 結果是靜默少算一條專家分支 —— 那是最嚴重的錯誤。 | |
| // 2000 次 × 250 µs = 最多等 0.5 秒;真的撐不到那麼久,代表 arena 被 | |
| // 預取的 busy 槽位佔住了(實測:跨層預取開啟時會發生)。 | |
| for (int attempt = 0; attempt < 2000 && s < 0; attempt++) { | |
| lk.unlock(); | |
| std::this_thread::sleep_for(std::chrono::microseconds(250)); | |
| lk.lock(); | |
| s = admit_slot_locked(il, ie); | |
| } | |
| if (s < 0) { | |
| fprintf(stderr, | |
| "[sdq] 嚴重:找不到能放 il=%d ie=%d(need_big=%d)的槽位," | |
| "大槽 %zu/小槽 %zu\n", | |
| il, ie, (int) need_big_slot(il, ie), n_big_, n_small_); | |
| // 這裡「回 -1」會讓呼叫端**跳過這個 expert**,結果是靜默少算一條專家分支 | |
| // —— 那是整個引擎最嚴重的錯誤型態(模型會安靜地給出錯誤答案,而且 | |
| // 每次執行還不一定一樣)。寧可整個行程死掉,也不要產出「看起來正常」 | |
| // 的錯誤結果。 | |
| fprintf(stderr, | |
| "[sdq] 致命:無法取得 expert 槽位。為了不給出靜默錯誤的結果," | |
| "直接中止。請回報這個情況(多半是預取把槽位佔成 busy 太久)。\n"); | |
| std::fflush(stderr); | |
| abort(); | |
| } | |
| } | |
| slot_state & st = slots_[s]; | |
| if (st.layer >= 0) { | |
| page_.erase(((int64_t) st.layer << 20) | (int64_t) st.expert); | |
| stats_.evictions.fetch_add(1, std::memory_order_relaxed); | |
| used_.fetch_sub(1, std::memory_order_relaxed); | |
| } | |
| st.layer = il; | |
| st.expert = ie; | |
| st.last_use = ++clock_; | |
| st.pin = 1; | |
| st.is_hot = hot_lookup(il, ie); | |
| // **先標 busy,再 publish 到 page_。** | |
| // 順序反過來的話(先 publish 再標 busy 中間有一個窗口)另一條執行緒會把它 | |
| // 當成命中,然後在資料還沒讀完時就拿去做矩陣乘法 → 非決定性的錯誤輸出。 | |
| // acquire() 看到 busy 會等;admit/victim 本來就會跳過 busy 的槽位。 | |
| st.busy = true; | |
| page_[((int64_t) il << 20) | (int64_t) ie] = s; | |
| lk.unlock(); | |
| const uint64_t t0 = now_us(); | |
| const bool ok = read_expert_par(il, ie, slot_ptr(s)); | |
| const uint64_t dt = now_us() - t0; | |
| std::unique_lock<std::mutex> lk2(mu_); | |
| st.busy = false; | |
| if (!ok) { | |
| st.layer = -1; | |
| st.expert = -1; | |
| st.pin = 0; | |
| page_.erase(((int64_t) il << 20) | (int64_t) ie); | |
| return -1; | |
| } | |
| used_.fetch_add(1, std::memory_order_relaxed); | |
| stats_.demand_wait_us.fetch_add(dt, std::memory_order_relaxed); | |
| const layer_desc * L = layer(il); | |
| if (L) { | |
| stats_.ssb_read_bytes_demand.fetch_add((uint64_t) L->experts[(size_t) ie].bytes, std::memory_order_relaxed); | |
| } | |
| return s; | |
| } | |
| void pager::prefetch_expert(int il, int ie) { | |
| if (!ready_ || !cfg_.prefetch) { | |
| return; | |
| } | |
| { | |
| std::unique_lock<std::mutex> lk(mu_); | |
| if (find_slot_locked(il, ie) >= 0) { | |
| stats_.prefetch_duplicate.fetch_add(1, std::memory_order_relaxed); | |
| return; | |
| } | |
| // 預取只能占用「空槽」或「已經很久沒被用到的頁」。 | |
| // (實測:無差別淘汰會讓命中率從 90% 掉到 27%,而預取使用率是 0%,純粹有害。) | |
| // | |
| // **槽位一定要選對尺寸**:need_big_slot() 為真的 expert 拿到小槽位會寫出邊界。 | |
| const uint64_t stale_before = clock_ - (uint64_t) env_size("SDQ_PREFETCH_STALE_TOKENS", 2) * 320; | |
| // 預取佇列也不能太深:SSD 讀不夠快時,排在後面的預取早就過期了 | |
| if (pf_q_.size() > (size_t) env_size("SDQ_PREFETCH_QUEUE", 64)) { | |
| return; | |
| } | |
| // **必須留夠空的槽位給 demand。** | |
| // 預取會把槽位標成 busy 並佔著它讀取;大槽位本來就少(工作集下限是 | |
| // n_big_layers × n_expert_used),如果預取把最後幾個空槽都吃掉, | |
| // 下一層的 demand 就會「找不到能放 il=34 ie=114 的槽位」而中止。 | |
| // (這是 verify 新增的預取壓力測試 [2b] 抓到的實際情況。) | |
| // 所以:對應尺寸的空槽少於 n_expert_used 個時,這一筆預取直接不做。 | |
| const size_t reserve = (size_t) info_.n_expert_used; | |
| const size_t n_free_right = need_big_slot(il, ie) ? free_big_slots_.size() : free_slots_.size(); | |
| if (n_free_right < reserve) { | |
| stats_.prefetch_duplicate.fetch_add(1, std::memory_order_relaxed); | |
| return; | |
| } | |
| const int32_t s = victim_stale_for(il, ie, stale_before); | |
| if (s < 0) { | |
| return; | |
| } | |
| slot_state & st = slots_[s]; | |
| if (st.layer >= 0) { | |
| page_.erase(((int64_t) st.layer << 20) | (int64_t) st.expert); | |
| stats_.evictions.fetch_add(1, std::memory_order_relaxed); | |
| used_.fetch_sub(1, std::memory_order_relaxed); | |
| } | |
| st.layer = il; | |
| st.expert = ie; | |
| st.busy = true; | |
| st.pin = 0; | |
| pf_q_.push_back({s, il, ie, now_us()}); | |
| stats_.prefetch_issued.fetch_add(1, std::memory_order_relaxed); | |
| } | |
| pf_cv_.notify_one(); | |
| } | |
| void pager::pf_worker() { | |
| std::unique_lock<std::mutex> lk(pf_mu_); | |
| for (;;) { | |
| pf_cv_.wait(lk, [this] { return pf_stop_ || !pf_q_.empty(); }); | |
| if (pf_stop_ && pf_q_.empty()) { | |
| return; | |
| } | |
| const pf_item item = pf_q_.front(); | |
| pf_q_.pop_front(); | |
| lk.unlock(); | |
| // 這個槽可能已經被淘汰/重用 → 檢查 generation | |
| { | |
| std::unique_lock<std::mutex> g(mu_); | |
| const slot_state & st = slots_[item.slot]; | |
| const bool still = st.busy && st.layer == item.il && st.expert == item.ie; | |
| if (!still) { | |
| lk.lock(); | |
| continue; | |
| } | |
| } | |
| // 預取也走 I/O 段平行:單獨一條執行緒循序讀三段,是這條路徑唯一的瓶頸。 | |
| const bool ok = cfg_.io_parallel ? read_expert_par(item.il, item.ie, slot_ptr(item.slot)) | |
| : read_expert(item.il, item.ie, slot_ptr(item.slot)); | |
| const layer_desc * L = layer(item.il); | |
| { | |
| std::unique_lock<std::mutex> g(mu_); | |
| slot_state & st = slots_[item.slot]; | |
| if (st.busy && st.layer == item.il && st.expert == item.ie) { | |
| st.busy = false; | |
| st.last_use = ++clock_; | |
| if (ok) { | |
| page_[((int64_t) item.il << 20) | (int64_t) item.ie] = item.slot; | |
| used_.fetch_add(1, std::memory_order_relaxed); | |
| stats_.prefetch_completed.fetch_add(1, std::memory_order_relaxed); | |
| if (L) { | |
| stats_.ssb_read_bytes_prefetch.fetch_add( | |
| (uint64_t) L->experts[(size_t) item.ie].bytes, std::memory_order_relaxed); | |
| } | |
| } else { | |
| st.layer = -1; | |
| st.expert = -1; | |
| } | |
| } | |
| } | |
| lk.lock(); | |
| } | |
| } | |
| void pager::release(int s) { | |
| if (s < 0 || (size_t) s >= n_slots_) { | |
| return; | |
| } | |
| std::lock_guard<std::mutex> lk(mu_); | |
| slot_state & st = slots_[(size_t) s]; | |
| if (st.pin > 0) { | |
| st.pin--; | |
| } | |
| } | |
| void pager::unpin_all() { | |
| std::unique_lock<std::mutex> lk(mu_); | |
| for (auto & st : slots_) { | |
| st.pin = 0; | |
| } | |
| } | |
| void pager::note_used(int il, int ie) { | |
| if (il < 0 || ie < 0) { | |
| return; | |
| } | |
| std::lock_guard<std::mutex> lk(hist_mu_); | |
| if (il >= (int) cur_used_.size()) { | |
| return; | |
| } | |
| auto & v = cur_used_[(size_t) il]; | |
| if (std::find(v.begin(), v.end(), ie) == v.end()) { | |
| v.push_back(ie); | |
| } | |
| } | |
| // 「下一層」預取:只在 layer L 的 op 開頭呼叫一次。 | |
| // | |
| // 為什麼需要這個函式:end_token() 是**整個 token 的 40 層都算完之後**才被呼叫的 | |
| // (sdq_cli.cpp 在 llama_decode 之後呼叫),所以它排的預取完全沒有和任何計算重疊 | |
| // —— 量測上 prefetch 使用率只有 0~8 %,等於整條預取路徑是關的。 | |
| // 在 layer L 開始時就替 layer L+1 排預取,才有一整層的時間(讀+算,約 4 ms) | |
| // 可以讓 SSD 讀取躲在後面。 | |
| // | |
| // 預測用的是 prev_used_[il+1](上一個 token 在該層選的專家)—— MoE 的路由在 | |
| // 鄰近 token 之間非常穩定,這是最便宜且不用訓練的預測器。 | |
| void pager::prefetch_layer_ahead(int il) { | |
| if (!ready_ || !cfg_.prefetch || il < 0 || il >= (int) layers_.size()) { | |
| return; | |
| } | |
| int32_t ids[64]; | |
| const int k = (int) info_.n_expert_used; | |
| if (k <= 0 || k > 64) { | |
| return; | |
| } | |
| predict(il, ids, k); | |
| for (int i = 0; i < k; i++) { | |
| prefetch_expert(il, ids[i]); | |
| } | |
| } | |
| void pager::end_token() { | |
| if (!ready_) { | |
| return; | |
| } | |
| // 用「上個 token 的選擇」更新一階 Markov 轉移表 | |
| { | |
| std::lock_guard<std::mutex> lk(hist_mu_); | |
| for (size_t il = 0; il < cur_used_.size(); il++) { | |
| auto & cur = cur_used_[il]; | |
| if (cur.empty()) { | |
| continue; | |
| } | |
| auto & layer_hist = hist_[il]; | |
| // 以每層最近一次(或平均)的前一次選擇為前驅 | |
| if (!prev_used_.empty() && il < prev_used_.size() && !prev_used_[il].empty()) { | |
| for (int32_t from : prev_used_[il]) { | |
| trans & t = layer_hist[(size_t) from]; | |
| for (int32_t to : cur) { | |
| int found = -1; | |
| for (int32_t i = 0; i < t.n; i++) { | |
| if (t.to[i] == to) { | |
| t.cnt[i]++; | |
| found = i; | |
| break; | |
| } | |
| } | |
| if (found < 0 && t.n < 64) { | |
| t.to[t.n] = to; | |
| t.cnt[t.n] = 1; | |
| t.n++; | |
| } | |
| } | |
| } | |
| } | |
| prev_used_[il] = cur; | |
| cur.clear(); | |
| } | |
| } | |
| // 依預測排預取 | |
| if (!cfg_.prefetch) { | |
| return; | |
| } | |
| int32_t ids[64]; | |
| for (size_t il = 0; il < layers_.size(); il++) { | |
| const int k = (int) info_.n_expert_used; | |
| predict((int) il, ids, k); | |
| for (int i = 0; i < k; i++) { | |
| prefetch_expert((int) il, ids[i]); | |
| } | |
| } | |
| } | |
| void pager::predict(int il, int32_t * out, int k) { | |
| if (il < 0 || il >= (int) prev_used_.size()) { | |
| for (int i = 0; i < k; i++) { | |
| out[i] = i; | |
| } | |
| return; | |
| } | |
| // 以最近一次選擇為前驅;沒有歷史就直接沿用上一次的選擇 | |
| const std::vector<int32_t> & prev = prev_used_[(size_t) il]; | |
| std::vector<std::pair<uint32_t, int32_t>> scored; | |
| if (prev.empty()) { | |
| for (int i = 0; i < k; i++) { | |
| out[i] = i; | |
| } | |
| return; | |
| } | |
| for (int32_t from : prev) { | |
| const trans & t = hist_[(size_t) il][(size_t) from]; | |
| for (int32_t i = 0; i < t.n; i++) { | |
| scored.push_back({t.cnt[i], t.to[i]}); | |
| } | |
| } | |
| if (scored.empty()) { | |
| for (size_t i = 0; i < prev.size() && i < (size_t) k; i++) { | |
| out[i] = prev[i]; | |
| } | |
| for (int i = (int) prev.size(); i < k; i++) { | |
| out[i] = i; | |
| } | |
| return; | |
| } | |
| std::sort(scored.begin(), scored.end(), [](auto & a, auto & b) { | |
| return a.first > b.first; | |
| }); | |
| for (int i = 0; i < k && i < (int) scored.size(); i++) { | |
| out[i] = scored[(size_t) i].second; | |
| } | |
| } | |
| double pager::hit_rate() const { | |
| const uint64_t req = stats_.expert_requests.load(); | |
| const uint64_t hit = stats_.expert_hits.load(); | |
| return req ? (double) hit / (double) req : 0.0; | |
| } | |
| void pager::shutdown() { | |
| if (!ready_) { | |
| return; | |
| } | |
| if (!profile_out_.empty()) { | |
| std::lock_guard<std::mutex> lk(mu_); | |
| dump_profile(profile_out_); | |
| } | |
| if (pf_.joinable()) { | |
| { | |
| std::unique_lock<std::mutex> lk(pf_mu_); | |
| pf_stop_ = true; | |
| } | |
| pf_cv_.notify_all(); | |
| pf_.join(); | |
| } | |
| io_pool_stop(); | |
| if (fd_ >= 0) { | |
| close(fd_); | |
| fd_ = -1; | |
| } | |
| if (fd_direct_ >= 0) { | |
| close(fd_direct_); | |
| fd_direct_ = -1; | |
| } | |
| ready_ = false; | |
| } | |
| pager::~pager() { | |
| shutdown(); | |
| if (arena_) { | |
| // 注意:arena 是 posix_memalign 配置的,必須用 free() 釋放, | |
| // 用 munmap() 會破壞 glibc 的 heap(之後任何 free 都會報 double free)。 | |
| free(arena_); | |
| arena_ = nullptr; | |
| } | |
| } | |
| static double pct(uint64_t a, uint64_t b) { | |
| return b ? 100.0 * (double) a / (double) b : 0.0; | |
| } | |
| void pager::print_summary() const { | |
| if (!ready_) { | |
| return; | |
| } | |
| const layer_desc * L0 = layer(0); | |
| if (L0 && !L0->experts.empty()) { | |
| const expert_desc & e = L0->experts[0]; | |
| fprintf(stderr, "[sdq] page-map: n_layer=%lld n_expert=%lld n_expert_used=%lld n_embd=%lld n_ff_exp=%lld" | |
| " | L0E0 gate=%llu/%llu up=%llu/%llu down=%llu/%llu (off/size) ne0=%lld ne1=%lld d_ne0=%lld d_ne1=%lld\n", | |
| (long long) info_.n_layer, (long long) info_.n_expert, (long long) info_.n_expert_used, | |
| (long long) info_.n_embd, (long long) info_.n_ff_exp, | |
| (unsigned long long) e.seg[part::GATE].off, (unsigned long long) e.seg[part::GATE].size, | |
| (unsigned long long) e.seg[part::UP].off, (unsigned long long) e.seg[part::UP].size, | |
| (unsigned long long) e.seg[part::DOWN].off, (unsigned long long) e.seg[part::DOWN].size, | |
| (long long) e.ne0, (long long) e.ne1, (long long) e.down_ne0, (long long) e.down_ne1); | |
| } | |
| fprintf(stderr, | |
| "[sdq] expert pager ready: %s\n" | |
| "[sdq] 熱門 expert 釘選:%d 個%s\n" | |
| "[sdq] slot=%.2f MiB (%zu B) slots=%zu arena=%.2f GiB (RAM 實測 RSS=%.2f GiB)\n", | |
| cfg_.model_path.c_str(), | |
| n_hot_, n_hot_ > 0 ? "(永不淘汰)" : "(無剖面,用動態 LRU/LFU)", | |
| slot_bytes_ / 1048576.0, slot_bytes_, n_slots_, arena_bytes_ / 1073741824.0, | |
| process_rss_bytes() / 1073741824.0); | |
| if (n_big_) { | |
| // 大小兩種槽位:把配置結果講清楚,否則只看「slots=N」會誤以為每個槽一樣大。 | |
| fprintf(stderr, | |
| "[sdq] 槽位分兩種:小 %.2f MiB × %zu + 大 %.2f MiB × %zu(平均 %.2f MiB)\n", | |
| small_bytes_ / 1048576.0, n_small_, | |
| big_bytes_ / 1048576.0, n_big_, | |
| (double) arena_bytes_ / (double) (n_slots_ ? n_slots_ : 1) / 1048576.0); | |
| } | |
| fprintf(stderr, | |
| "[sdq] prefetch=%s fadvise_dontneed=%d\n", | |
| cfg_.prefetch ? "on" : "off", (int) cfg_.fadvise_dontneed); | |
| } | |
| void pager::dump_stats(const char * phase) { | |
| if (cfg_.stats_path.empty()) { | |
| return; | |
| } | |
| const uint64_t tokens = stats_.tokens.load(); | |
| const uint64_t req = stats_.expert_requests.load(); | |
| const uint64_t hit = stats_.expert_hits.load(); | |
| const uint64_t issued = stats_.prefetch_issued.load(); | |
| const uint64_t used = stats_.prefetch_used.load(); | |
| const uint64_t b_dem = stats_.ssb_read_bytes_demand.load(); | |
| const uint64_t b_pre = stats_.ssb_read_bytes_prefetch.load(); | |
| std::ostringstream o; | |
| o << "{\"phase\":\"" << phase << "\"" | |
| << ",\"tokens\":" << tokens | |
| << ",\"expert_requests\":" << req | |
| << ",\"expert_hits\":" << hit | |
| << ",\"expert_misses\":" << (req - hit) | |
| << ",\"hit_rate\":" << pct(hit, req) | |
| << ",\"ssb_demand_bytes\":" << b_dem | |
| << ",\"ssb_prefetch_bytes\":" << b_pre | |
| << ",\"ssb_total_bytes\":" << (b_dem + b_pre) | |
| << ",\"ssd_bytes_per_token\":" << (tokens ? (double) (b_dem + b_pre) / (double) tokens : 0.0) | |
| << ",\"prefetch_issued\":" << issued | |
| << ",\"prefetch_used\":" << used | |
| << ",\"prefetch_used_rate\":" << pct(used, issued) | |
| << ",\"prefetch_duplicate\":" << stats_.prefetch_duplicate.load() | |
| << ",\"evictions\":" << stats_.evictions.load() | |
| << ",\"demand_wait_us\":" << stats_.demand_wait_us.load() | |
| << ",\"prefetch_wait_us\":" << stats_.prefetch_wait_us.load() | |
| << ",\"arena_bytes\":" << arena_bytes_ | |
| // 大小兩種槽位之後,「用掉幾個槽位」不能直接乘 slot_bytes_(那是大的那種)。 | |
| // dump_stats 只在階段轉換時呼叫一次,遍歷兩千幾個槽位的成本可以忽略。 | |
| << ",\"arena_used_bytes\":" << used_bytes_snapshot() | |
| << ",\"rss_bytes\":" << process_rss_bytes() | |
| << ",\"n_slots\":" << n_slots_ | |
| << ",\"n_small_slots\":" << n_small_ | |
| << ",\"n_big_slots\":" << n_big_ | |
| << ",\"small_slot_bytes\":" << small_bytes_ | |
| << ",\"big_slot_bytes\":" << big_bytes_ | |
| << ",\"slot_bytes\":" << slot_bytes_ | |
| << "}"; | |
| std::ofstream f(cfg_.stats_path, std::ios::app); | |
| if (f) { | |
| f << o.str() << "\n"; | |
| } | |
| } | |
| } // namespace sdq | |
| namespace sdq { | |
| // 由 CLI 呼叫:依環境變數/參數初始化 pager | |
| bool sdq_pager_active_for_layer(int il) { | |
| pager & P = pager::instance(); | |
| if (!P.ready() || P.layer(il) == nullptr) { | |
| return false; | |
| } | |
| // 除錯用:SDQ_MOE_LO / SDQ_MOE_HI 可以只讓部分層走分頁路徑, | |
| // 用來二分搜尋第一個發生數值差異的層。 | |
| static const int lo = [] { const char * e = getenv("SDQ_MOE_LO"); return e ? atoi(e) : 0; }(); | |
| static const int hi = [] { const char * e = getenv("SDQ_MOE_HI"); return e ? atoi(e) : 1 << 30; }(); | |
| return il >= lo && il < hi; | |
| } | |
| bool sdq_init(const std::string & model_path) { | |
| config c; | |
| c.model_path = model_path; | |
| c.arena_bytes = env_size("SDQ_ARENA_MB", 0) << 20; | |
| c.ram_budget_bytes = env_size("SDQ_RAM_BUDGET_MB", 0) << 20; | |
| c.slot_bytes = env_size("SDQ_SLOT_KB", 0) << 10; | |
| c.prefetch = env_flag("SDQ_PREFETCH", true); | |
| // ---- OPT-1:熱門 expert 剖面(Strata 式)---- | |
| if (const char * pi = getenv("SDQ_PROFILE_IN")) { c.profile_in = pi; } | |
| if (const char * po = getenv("SDQ_PROFILE_OUT")) { c.profile_out = po; } | |
| c.hot_pin = (int) env_size("SDQ_HOT_PIN", 0); | |
| c.admit_min_hits = (int) env_size("SDQ_ADMIT_MIN_HITS", 1); | |
| c.fadvise_dontneed = env_flag("SDQ_FADV_DONTNEED", true); | |
| c.io_parallel = env_flag("SDQ_IO_PARALLEL", true); | |
| c.io_threads = (int) env_size("SDQ_IO_THREADS", 0); | |
| c.io_split = (int) env_size("SDQ_IO_SPLIT", 1); | |
| c.verbose = env_flag("SDQ_VERBOSE", false); | |
| const char * sp = getenv("SDQ_STATS_FILE"); | |
| if (sp) { | |
| c.stats_path = sp; | |
| } | |
| sdq_set_config(c); | |
| try { | |
| pager::instance().init(c, nullptr); | |
| } catch (const std::exception & e) { | |
| fprintf(stderr, "[sdq] pager 初始化失敗:%s\n", e.what()); | |
| return false; | |
| } | |
| return pager::instance().ready(); | |
| } | |
| // ---------------------------------------------------------------- 除錯:dump 任意 F32 張量 | |
| // 目的:讓「原生 mul_mat_id 路徑」與「SDQ 分頁路徑」輸出同一層的 MoE 結果, | |
| // 兩邊各自 dump 一次就能直接 A/B 對照,判定到底是哪一邊算錯。 | |
| namespace { | |
| struct dump_arg { std::string path; }; | |
| void dump_op(ggml_tensor * dst, int ith, int /*nth*/, void * userdata) { | |
| if (ith != 0) { | |
| return; | |
| } | |
| const ggml_tensor * src = dst->src[0]; | |
| dump_arg * a = (dump_arg *) userdata; | |
| FILE * f = fopen(a->path.c_str(), "wb"); | |
| if (!f) { | |
| fprintf(stderr, "[sdq-dump] 無法開啟 %s\n", a->path.c_str()); | |
| return; | |
| } | |
| const int64_t meta[6] = { src->ne[0], src->ne[1], src->ne[2], src->ne[3], | |
| (int64_t) ggml_nbytes(src), 0 }; | |
| fwrite(meta, sizeof(int64_t), 6, f); | |
| // 只支援 F32 contiguous;其他情況就明確報錯,不要假設 | |
| if (src->type != GGML_TYPE_F32 || !ggml_is_contiguous(src)) { | |
| fprintf(stderr, "[sdq-dump] %s 不是 contiguous F32(type=%d)\n", | |
| a->path.c_str(), (int) src->type); | |
| } else { | |
| fwrite(src->data, sizeof(float), (size_t) (ggml_nelements(src)), f); | |
| } | |
| fclose(f); | |
| fprintf(stderr, "[sdq-dump] %s (%lld 個元素)\n", a->path.c_str(), | |
| (long long) ggml_nelements(src)); | |
| } | |
| } // namespace | |
| ggml_tensor * sdq_maybe_dump_tensor(ggml_context * ctx, ggml_tensor * t, int il) { | |
| const char * prefix = getenv("SDQ_MOE_OUT_DUMP"); | |
| if (!prefix || !t) { | |
| return t; | |
| } | |
| int layer = getenv("SDQ_MOE_OUT_DUMP_LAYER") ? atoi(getenv("SDQ_MOE_OUT_DUMP_LAYER")) : -1; | |
| if (layer >= 0 && il != layer) { | |
| return t; | |
| } | |
| dump_arg * a = new dump_arg{ std::string(prefix) + "." + std::to_string(il) }; | |
| return ggml_custom_4d(ctx, GGML_TYPE_F32, t->ne[0], t->ne[1], t->ne[2], t->ne[3], | |
| &t, 1, dump_op, 1, a); | |
| } | |
| void sdq_shutdown() { | |
| pager & P = pager::instance(); | |
| sdq_print_time_profile(); | |
| P.dump_stats("shutdown"); | |
| P.print_summary(); | |
| P.shutdown(); | |
| } | |
| } // namespace sdq | |