Skip to content

Commit 2998033

Browse files
Add hf3fs support for hicache storage (based on #7704) (#7280)
Co-authored-by: Zhiqiang Xie <xiezhq@stanford.edu>
1 parent a79a5d7 commit 2998033

12 files changed

Lines changed: 1110 additions & 23 deletions

File tree

benchmark/hf3fs/bench.sh

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
1+
SGLANG_HICACHE_HF3FS_CONFIG_PATH=/sgl-workspace/sglang/benchmark/hf3fs/hf3fs.json \
2+
python3 benchmark/hf3fs/bench_storage.py
3+
4+
####################################################################################################
5+
6+
rm -rf nohup.out && \
7+
nohup python3 -m sglang.launch_server \
8+
--model-path /code/models/Qwen3-32B/ \
9+
--host 0.0.0.0 --port 33301 \
10+
--page-size 64 \
11+
--enable-hierarchical-cache \
12+
--hicache-ratio 2 --hicache-size 0 \
13+
--hicache-write-policy write_through \
14+
--hicache-storage-backend hf3fs &
15+
16+
rm -rf bench_multiturn.out && \
17+
nohup python3 benchmark/hicache/bench_multiturn.py \
18+
--model-path /code/models/Qwen3-32B \
19+
--dataset-path /code/models/ShareGPT_V3_unfiltered_cleaned_split.json \
20+
--port 33301 \
21+
--request-length 2048 --num-clients 512 --num-rounds 3 --max-parallel 8 \
22+
> bench_multiturn.out &
23+
24+
####################################################################################################
25+
26+
rm -rf nohup.out && \
27+
nohup python3 -m sglang.launch_server \
28+
--model-path /code/models/DeepSeek-R1/ \
29+
--tp 16 --nnodes 2 --node-rank 0 \
30+
--dist-init-addr 10.74.249.153:5000 \
31+
--host 0.0.0.0 --port 33301 \
32+
--page-size 64 \
33+
--enable-hierarchical-cache \
34+
--hicache-ratio 2 --hicache-size 60 \
35+
--hicache-write-policy write_through \
36+
--hicache-storage-backend hf3fs &
37+
38+
rm -rf bench_multiturn.out && \
39+
nohup python3 benchmark/hicache/bench_multiturn.py \
40+
--model-path /code/models/Qwen3-32B \
41+
--dataset-path /code/models/ShareGPT_V3_unfiltered_cleaned_split.json \
42+
--port 33301 \
43+
--request-length 2048 --num-clients 1024 --num-rounds 3 --max-parallel 8 \
44+
> bench_multiturn.out &
45+
46+
####################################################################################################
47+
48+
ps aux | grep "sglang.launch_server" | grep -v grep | awk '{print $2}' | xargs kill -9
49+
ps aux | grep "bench_multiturn.py" | grep -v grep | awk '{print $2}' | xargs kill -9

benchmark/hf3fs/bench_client.py

Lines changed: 162 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,162 @@
1+
import concurrent.futures
2+
import logging
3+
import random
4+
import time
5+
from typing import List
6+
7+
import torch
8+
from tqdm import tqdm
9+
10+
from sglang.srt.mem_cache.storage.hf3fs.client_hf3fs import Hf3fsClient
11+
12+
13+
def print_stats(x: List[int]):
14+
x = sorted(x)
15+
lenx = len(x)
16+
print(
17+
f"mean = {sum(x)/len(x):.2f}, "
18+
f"min = {min(x):.2f}, "
19+
f"p25 = {x[int(lenx*0.25)]:.2f}, "
20+
f"p50 = {x[int(lenx*0.5)]:.2f}, "
21+
f"p75 = {x[int(lenx*0.75)]:.2f}, "
22+
f"max = {max(x):.2f}"
23+
)
24+
25+
26+
def test():
27+
# /path/to/hf3fs
28+
file_path = "/data/bench.bin"
29+
file_size = 1 << 40
30+
bytes_per_page = 16 << 20
31+
entries = 32
32+
file_ops = Hf3fsClient(file_path, file_size, bytes_per_page, entries)
33+
34+
print("test batch_read / batch_write")
35+
num_pages = 128
36+
dtype = torch.bfloat16
37+
numel = bytes_per_page // dtype.itemsize
38+
offsets = list(range(file_size // bytes_per_page))
39+
random.shuffle(offsets)
40+
offsets = offsets[:num_pages]
41+
offsets = [i * bytes_per_page for i in offsets]
42+
tensor_writes = [
43+
torch.randn(numel, dtype=dtype)
44+
for _ in tqdm(range(num_pages), desc="prepare tensor")
45+
]
46+
for i in tqdm(range(0, num_pages, file_ops.entries), desc="batch_write"):
47+
results = file_ops.batch_write(
48+
offsets[i : i + file_ops.entries], tensor_writes[i : i + file_ops.entries]
49+
)
50+
assert all([result == numel * dtype.itemsize for result in results])
51+
tensor_reads = [
52+
torch.empty(numel, dtype=dtype)
53+
for _ in tqdm(range(num_pages), desc="prepare tensor")
54+
]
55+
for i in tqdm(range(0, num_pages, file_ops.entries), desc="batch_read"):
56+
results = file_ops.batch_read(
57+
offsets[i : i + file_ops.entries], tensor_reads[i : i + file_ops.entries]
58+
)
59+
assert all([result == numel * dtype.itemsize for result in results])
60+
assert all([torch.allclose(r, w) for r, w in zip(tensor_reads, tensor_writes)])
61+
62+
file_ops.close()
63+
print("test done")
64+
65+
66+
def bench():
67+
file_path = "/data/bench.bin"
68+
file_size = 1 << 40
69+
bytes_per_page = 16 << 20
70+
entries = 8
71+
numjobs = 16
72+
73+
dtype = torch.bfloat16
74+
numel = bytes_per_page // dtype.itemsize
75+
76+
file_ops = [
77+
Hf3fsClient(file_path, file_size, bytes_per_page, entries)
78+
for _ in range(numjobs)
79+
]
80+
81+
num_page = entries
82+
83+
offsets = list(range(file_size // bytes_per_page))
84+
tensors_write = [torch.randn(numel, dtype=dtype)] * num_page
85+
tensors_read = [torch.empty(numel, dtype=dtype)] * num_page
86+
random.shuffle(offsets)
87+
88+
warmup = 50
89+
iteration = 100
90+
91+
executor = concurrent.futures.ThreadPoolExecutor(max_workers=numjobs)
92+
93+
w_bw = []
94+
w_size = num_page * numjobs * bytes_per_page / (1 << 30)
95+
for i in tqdm(range(warmup + iteration), desc="Benchmarking write (GB/s)"):
96+
_offsets = [
97+
[
98+
offset * bytes_per_page
99+
for offset in offsets[
100+
(i * numjobs + j) * num_page : (i * numjobs + j + 1) * num_page
101+
]
102+
]
103+
for j in range(numjobs)
104+
]
105+
tik = time.perf_counter()
106+
futures = [
107+
executor.submit(file_ops[j].batch_write, offset, tensors_write)
108+
for j, offset in enumerate(_offsets)
109+
]
110+
results = [future.result() for future in futures]
111+
tok = time.perf_counter()
112+
if i < warmup:
113+
continue
114+
w_bw.append(w_size / (tok - tik))
115+
results = [
116+
_result == bytes_per_page for result in results for _result in result
117+
]
118+
assert all(results)
119+
print_stats(w_bw)
120+
121+
r_bw = []
122+
r_size = w_size
123+
for i in tqdm(range(warmup + iteration), desc="Benchmarking read (GB/s)"):
124+
_offsets = [
125+
[
126+
offset * bytes_per_page
127+
for offset in offsets[
128+
(i * numjobs + j) * num_page : (i * numjobs + j + 1) * num_page
129+
]
130+
]
131+
for j in range(numjobs)
132+
]
133+
tik = time.perf_counter()
134+
futures = [
135+
executor.submit(file_ops[j].batch_read, offset, tensors_read)
136+
for j, offset in enumerate(_offsets)
137+
]
138+
results = [future.result() for future in futures]
139+
tok = time.perf_counter()
140+
if i < warmup:
141+
continue
142+
r_bw.append(r_size / (tok - tik))
143+
results = [
144+
_result == bytes_per_page for result in results for _result in result
145+
]
146+
assert all(results)
147+
print_stats(r_bw)
148+
149+
executor.shutdown(wait=True)
150+
for _file_ops in file_ops:
151+
_file_ops.close()
152+
print("bench done")
153+
154+
155+
def main():
156+
logging.basicConfig(level=logging.INFO)
157+
test()
158+
bench()
159+
160+
161+
if __name__ == "__main__":
162+
main()

0 commit comments

Comments
 (0)