From f55a59f83eaaa4087f57e57d12227596d09af21f Mon Sep 17 00:00:00 2001 From: Feng Ren Date: Tue, 17 Mar 2026 10:43:23 +0800 Subject: [PATCH] [Bench] SSD benchmarks based on Mooncake Trace (#1613) --- .../storage_benchmark/storage_benchmark.py | 956 ++++++++++++++++++ docs/source/performance/storage-benchmark.md | 163 +++ 2 files changed, 1119 insertions(+) create mode 100644 benchmarks/storage_benchmark/storage_benchmark.py create mode 100644 docs/source/performance/storage-benchmark.md diff --git a/benchmarks/storage_benchmark/storage_benchmark.py b/benchmarks/storage_benchmark/storage_benchmark.py new file mode 100644 index 00000000..689bc98e --- /dev/null +++ b/benchmarks/storage_benchmark/storage_benchmark.py @@ -0,0 +1,956 @@ +#!/usr/bin/env python3 +# SPDX-License-Identifier: Apache-2.0 + +""" +Mooncake KVCache Storage Benchmark Tool +""" + +import argparse +import json +import time +import os +import statistics +import random +import errno +from pathlib import Path +from typing import Dict, List, Optional +from dataclasses import dataclass + +# ============================================================================ +# Constants +# ============================================================================ + +BLOCK_SIZE_TOKENS = 512 # Number of tokens per block +DEFAULT_BYTES_PER_TOKEN = 2048 # 7B model FP16 (2KB per token) +BLOCK_SIZE_BYTES = BLOCK_SIZE_TOKENS * DEFAULT_BYTES_PER_TOKEN # 1MB per block +MIN_LATENCY_MS = 0.001 # Minimum latency in milliseconds (1 microsecond) + +# Model KVCache sizes (bytes per token, based on LMCache calculator) +# Source: https://lmcache.ai/kv_cache_calculator.html +MODEL_BYTES_PER_TOKEN = { + "llama-3.1-405b": 327680, + "qwen3-32b": 81920, + "deepseek-v3": 1748992, + "glm-4.6": 157013, + "default": DEFAULT_BYTES_PER_TOKEN, +} + +# ============================================================================ +# Data Structures +# ============================================================================ + +@dataclass +class KVCacheRequest: + """KVCache request + + Attributes: + timestamp: Request timestamp in milliseconds + hash_ids: List of block IDs (each ID corresponds to a 512-token block) + input_length: Input token count + output_length: Output token count + """ + timestamp: float + hash_ids: List[int] + input_length: int + output_length: int + +# ============================================================================ +# Storage Layer: Offset Allocator +# ============================================================================ + +class OffsetAllocatorStorage: + """High-performance block storage based on Offset Allocator + + Architecture: + ----------- + 1. Single large file stores all blocks (avoids file explosion) + 2. Uses offset to manage file space (similar to Mooncake's OffsetAllocator) + 3. hash_id -> offset mapping stored in memory (fast lookup) + + Block Organization: + ----------- + Each block corresponds to 512 tokens, fixed size 1MB: + - hash_id[0] -> block_0 (tokens [0...511]) -> offset 0 + - hash_id[1] -> block_1 (tokens [512...1023]) -> offset 1 + - hash_id[i] -> block_i (tokens [i*512...(i+1)*512-1]) -> offset i + + Performance Advantages: + ----------- + - Only one file, no file explosion + - Offset reuse, reduces memory allocation + - pread/pwrite, thread-safe, no seek needed + - Keep fd open, reduces open/close overhead + - Metadata in memory, O(1) lookup + + Attributes: + storage_dir: Storage directory path + block_size_bytes: Block size in bytes + max_blocks: Maximum number of blocks + hash_id_to_offset: hash_id -> offset mapping + free_offsets: List of reusable offsets + next_offset: Next allocatable offset + """ + + def __init__(self, storage_dir: str, bytes_per_token: int = DEFAULT_BYTES_PER_TOKEN, + max_blocks: int = 100000, block_size_tokens: int = 512, + fsync_mode: str = 'batch', fsync_batch_size: int = 100): + """Initialize Offset Allocator storage + + Args: + storage_dir: Storage directory path + bytes_per_token: Bytes per token + max_blocks: Maximum number of blocks (determines file size) + block_size_tokens: Number of tokens per block + fsync_mode: When to fsync ('batch', 'always', 'end', 'none') + fsync_batch_size: Number of writes between fsync in batch mode + """ + self.storage_dir = Path(storage_dir) + self.bytes_per_token = bytes_per_token + self.block_size_tokens = block_size_tokens + self.block_size_bytes = self.block_size_tokens * self.bytes_per_token + self.max_blocks = max_blocks + + # Fsync configuration + self.fsync_mode = fsync_mode + self.fsync_batch_size = fsync_batch_size + self.pending_sync_count = 0 + + # Create storage directory + self.storage_dir.mkdir(parents=True, exist_ok=True) + + # Single large file + self.storage_file = self.storage_dir / "kvcache_storage.bin" + self.file_size = self.max_blocks * self.block_size_bytes + + # Initialize storage file + if not self.storage_file.exists(): + self._init_storage_file() + + # hash_id -> offset mapping (metadata, in memory) + self.hash_id_to_offset: Dict[int, int] = {} + + # Offset allocator (free list) + self.free_offsets: List[int] = [] + self.next_offset = 0 + + # File descriptor (keep open, avoid repeated open/close) + self.fd = None + + # Pre-allocated data buffer with pattern to avoid SSD compression artifacts + # Using a repeating pattern that looks like realistic data (not all zeros) + # Pattern: 64-byte repeated sequence mixed with some variation + pattern = bytes([(i & 0xFF) for i in range(256)]) # 0-255 byte pattern + pattern_repeats = (self.block_size_bytes // len(pattern)) + 1 + self._data_buffer = (pattern * pattern_repeats)[:self.block_size_bytes] + + # Statistics + self.stats = { + 'read_count': 0, + 'write_count': 0, + 'read_bytes': 0, + 'write_bytes': 0, + 'read_latencies_ms': [], + 'write_latencies_ms': [], + 'sync_count': 0, # Number of fsync operations performed + } + + # ======================================================================== + # Internal Methods + # ======================================================================== + + def _init_storage_file(self): + """Initialize storage file (pre-allocate space) + + Create sparse file to avoid actual disk space usage until data is written + """ + with open(self.storage_file, 'wb') as f: + f.seek(self.file_size - 1) + f.write(b'\0') + f.flush() + os.fsync(f.fileno()) + + def _get_fd(self): + """Get file descriptor (lazy open) + + Returns: + int: File descriptor + """ + if self.fd is None: + # Use O_RDWR | O_CREAT, no O_DIRECT (Python compatibility) + self.fd = os.open(self.storage_file, os.O_RDWR | os.O_CREAT) + return self.fd + + def _allocate_offset(self) -> int: + """Allocate a new offset + + Prioritize reusing freed offsets, otherwise allocate new offset + + Returns: + int: Allocated offset + """ + if self.free_offsets: + return self.free_offsets.pop() + offset = self.next_offset + self.next_offset += 1 + return offset + + def _free_offset(self, offset: int): + """Free offset for reuse + + Args: + offset: Offset to free + """ + self.free_offsets.append(offset) + + # ======================================================================== + # Public Interface + # ======================================================================== + + def block_exists(self, hash_id: int) -> bool: + """Check if block exists + + Args: + hash_id: Unique block identifier + + Returns: + bool: Whether block exists + """ + return hash_id in self.hash_id_to_offset + + def read_block(self, hash_id: int) -> float: + """Read block using pread + + Args: + hash_id: Unique block identifier + + Returns: + float: Read latency in milliseconds, or 0 if block doesn't exist + """ + if hash_id not in self.hash_id_to_offset: + return 0.0 # Block doesn't exist, no latency to measure + + offset = self.hash_id_to_offset[hash_id] + file_offset = offset * self.block_size_bytes + + start = time.perf_counter() + + try: + fd = self._get_fd() + data = os.pread(fd, self.block_size_bytes, file_offset) + latency_ms = (time.perf_counter() - start) * 1000.0 + + self.stats['read_count'] += 1 + self.stats['read_bytes'] += len(data) + self.stats['read_latencies_ms'].append(latency_ms) + return latency_ms + except OSError as e: + print(f"Error reading block {hash_id} at offset {file_offset}: {e}") + return 0.0 # Error case, don't pollute stats + + def write_block(self, hash_id: int) -> float: + """Write block using pwrite + + Args: + hash_id: Unique block identifier + + Returns: + float: Write latency in milliseconds + """ + # Allocate offset + offset = self._allocate_offset() + file_offset = offset * self.block_size_bytes + + # Use pre-allocated buffer (much faster than os.urandom) + data = self._data_buffer + + start = time.perf_counter() + + try: + fd = self._get_fd() + written = os.pwrite(fd, data, file_offset) + + write_done = time.perf_counter() + + # Conditional fsync based on mode + if self.fsync_mode == 'always': + # Include fsync in latency measurement + os.fsync(fd) + self.stats['sync_count'] += 1 + self.pending_sync_count = 0 + latency_ms = (time.perf_counter() - start) * 1000.0 + # Evict from page cache AFTER fsync to ensure reads measure actual SSD performance + os.posix_fadvise(fd, file_offset, self.block_size_bytes, os.POSIX_FADV_DONTNEED) + elif self.fsync_mode == 'batch': + # For batch mode, only measure write time (fsync is deferred) + self.pending_sync_count += 1 + if self.pending_sync_count >= self.fsync_batch_size: + os.fsync(fd) + self.stats['sync_count'] += 1 + self.pending_sync_count = 0 + latency_ms = (write_done - start) * 1000.0 # Only write time + # Evict from page cache after each write + os.posix_fadvise(fd, file_offset, self.block_size_bytes, os.POSIX_FADV_DONTNEED) + elif self.fsync_mode == 'none': + latency_ms = (write_done - start) * 1000.0 + # Evict from page cache even when not syncing + os.posix_fadvise(fd, file_offset, self.block_size_bytes, os.POSIX_FADV_DONTNEED) + else: # 'end' mode + latency_ms = (write_done - start) * 1000.0 + # Evict from page cache (fsync will happen at the end) + os.posix_fadvise(fd, file_offset, self.block_size_bytes, os.POSIX_FADV_DONTNEED) + + # Update mapping + self.hash_id_to_offset[hash_id] = offset + + self.stats['write_count'] += 1 + self.stats['write_bytes'] += written + self.stats['write_latencies_ms'].append(latency_ms) + return latency_ms + except OSError as e: + if e.errno == errno.ENOSPC: + print(f"Error: Disk full when writing block {hash_id} at offset {file_offset}") + else: + print(f"Error writing block {hash_id} at offset {file_offset}: {e}") + return 0.0 # Error case, don't pollute stats + + def __enter__(self): + """Context manager entry""" + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + """Context manager exit - ensures cleanup""" + # Perform final fsync before closing for 'end' and 'batch' modes + self._finalize_sync() + self.close(force_sync=False) # Already synced above + return False + + def _finalize_sync(self): + """Perform final fsync before closing (for 'end' mode and pending batch writes)""" + if self.fd is not None: + if self.fsync_mode == 'end': + try: + os.fsync(self.fd) + self.stats['sync_count'] += 1 + except OSError: + pass + elif self.fsync_mode == 'batch' and self.pending_sync_count > 0: + # Flush remaining pending writes + try: + os.fsync(self.fd) + self.stats['sync_count'] += 1 + self.pending_sync_count = 0 + except OSError: + pass + + def close(self, force_sync: bool = True): + """Close file + + Args: + force_sync: Whether to force fsync before closing + """ + # For backward compatibility with non-context-manager usage + if force_sync: + self._finalize_sync() + + if self.fd is not None: + os.close(self.fd) + self.fd = None + + def get_stats(self) -> Dict: + """Get statistics + + Returns: + Dict: Dictionary containing read/write statistics + """ + def calc_stats(latencies): + """Calculate latency statistics""" + if not latencies: + return {'avg_ms': 0, 'p50_ms': 0, 'p95_ms': 0, 'p99_ms': 0} + return { + 'avg_ms': statistics.mean(latencies), + **calc_percentiles(latencies), + } + + return { + 'read': { + 'count': self.stats['read_count'], + 'mb': self.stats['read_bytes'] / 1024 / 1024, + **calc_stats(self.stats['read_latencies_ms']) + }, + 'write': { + 'count': self.stats['write_count'], + 'mb': self.stats['write_bytes'] / 1024 / 1024, + **calc_stats(self.stats['write_latencies_ms']) + }, + 'sync_count': self.stats['sync_count'], + 'total_blocks': len(self.hash_id_to_offset), + 'free_blocks': len(self.free_offsets), + } + + +# ============================================================================ +# Benchmark Layer +# ============================================================================ + +class StorageBenchmark: + """KVCache storage benchmark + + Based on Mooncake OffsetAllocator + vLLM PagedAttention implementation: + + Example: + ----- + Request A: [1, 2, 4] + -> hash_id 1 -> not exist, write block_1 (offset=0, 1MB) + -> hash_id 2 -> not exist, write block_2 (offset=1, 1MB) + -> hash_id 4 -> not exist, write block_4 (offset=2, 1MB) + + Request B: [1, 2, 4, 6] + -> hash_id 1 -> exists, read block_1 (offset=0) ✓ prefix reuse + -> hash_id 2 -> exists, read block_2 (offset=1) ✓ prefix reuse + -> hash_id 4 -> exists, read block_4 (offset=2) ✓ prefix reuse + -> hash_id 6 -> not exist, write block_6 (offset=3, 1MB) + + Performance Advantages: + --------- + - Single file operation, no file explosion + - Offset reuse, reduces memory allocation + - pread/pwrite, thread-safe + """ + + def __init__(self, storage_dir: str, bytes_per_token: int = DEFAULT_BYTES_PER_TOKEN, + max_blocks: int = 100000, block_size_tokens: int = 512, + fsync_mode: str = 'batch', fsync_batch_size: int = 100): + """Initialize benchmark + + Args: + storage_dir: Storage directory + bytes_per_token: Bytes per token + max_blocks: Maximum number of blocks + block_size_tokens: Number of tokens per block + fsync_mode: When to fsync ('batch', 'always', 'end', 'none') + fsync_batch_size: Number of writes between fsync in batch mode + """ + self.storage = OffsetAllocatorStorage( + storage_dir, bytes_per_token, max_blocks, + block_size_tokens, fsync_mode, fsync_batch_size + ) + self.bytes_per_token = bytes_per_token + self.block_size_tokens = block_size_tokens + + # Statistics + self.stats = { + 'total_requests': 0, + 'total_blocks': 0, + 'read_blocks': 0, + 'write_blocks': 0, + 'prefix_hit_blocks': 0, # Number of prefix hit blocks + 'request_latencies_ms': [], + } + + def process_request(self, req: KVCacheRequest) -> float: + """Process a KVCache request + + Based on vLLM's prefix caching mechanism: + - Each hash_id corresponds to an independent block + - Prefix reuse achieved through hash_id matching + + Args: + req: KVCache request + + Returns: + float: Request latency in milliseconds + """ + self.stats['total_requests'] += 1 + self.stats['total_blocks'] += len(req.hash_ids) + + start_time = time.perf_counter() + total_latency = 0.0 + + # Process each hash_id (in order) + for hash_id in req.hash_ids: + if self.storage.block_exists(hash_id): + # Block exists, read (reuse cached block) + total_latency += self.storage.read_block(hash_id) + self.stats['read_blocks'] += 1 + self.stats['prefix_hit_blocks'] += 1 # Count all cache hits as prefix reuse + else: + # Block doesn't exist, write (new block) + total_latency += self.storage.write_block(hash_id) + self.stats['write_blocks'] += 1 + + latency_ms = total_latency if total_latency > 0 else MIN_LATENCY_MS + self.stats['request_latencies_ms'].append(latency_ms) + + return latency_ms + + def get_stats(self) -> Dict: + """Get statistics + + Returns: + Dict: Statistics dictionary + """ + storage_stats = self.storage.get_stats() + + request_latencies = self.stats['request_latencies_ms'] + + if request_latencies: + latency_stats = { + 'avg_ms': statistics.mean(request_latencies), + **calc_percentiles(request_latencies), + } + else: + latency_stats = {'avg_ms': 0, 'p50_ms': 0, 'p95_ms': 0, 'p99_ms': 0} + + total_blocks = self.stats['total_blocks'] + read_blocks = self.stats['read_blocks'] + write_blocks = self.stats['write_blocks'] + + return { + 'total_requests': self.stats['total_requests'], + 'total_blocks': total_blocks, + 'read_blocks': read_blocks, + 'write_blocks': write_blocks, + 'prefix_hit_blocks': self.stats['prefix_hit_blocks'], + 'block_hit_rate': read_blocks / total_blocks if total_blocks > 0 else 0, + 'write_ratio': write_blocks / total_blocks if total_blocks > 0 else 0, + 'tokens_per_block': self.block_size_tokens, # Configurable block size in tokens + 'latency': latency_stats, + 'storage': storage_stats, + } + + def __enter__(self): + """Context manager entry""" + return self + + def __exit__(self, exc_type, exc_val, exc_tb): + """Context manager exit - ensures cleanup""" + self.close() + return False + + def close(self, force_sync: bool = True): + """Close storage + + Args: + force_sync: Whether to force final sync before closing + """ + self.storage.close(force_sync=force_sync) + + +# ============================================================================ +# Utility Functions +# ============================================================================ + +def calc_percentiles(data: List[float]) -> Dict[str, float]: + """Calculate latency percentiles + + Uses linear interpolation for accurate percentile calculation. + This is more accurate than statistics.quantiles() for small datasets. + + Args: + data: List of latency values in milliseconds + + Returns: + Dict containing p50, p95, p99 percentiles + """ + if not data: + return {'p50_ms': 0, 'p95_ms': 0, 'p99_ms': 0} + + # Sort data for percentile calculation + sorted_data = sorted(data) + n = len(sorted_data) + + def get_percentile(p: float) -> float: + """Get percentile using linear interpolation + + Args: + p: Percentile (0-100) + + Returns: + Value at percentile + """ + index = (n - 1) * p / 100 + lower = int(index) + upper = min(lower + 1, n - 1) + + if lower == upper: + return sorted_data[lower] + + # Linear interpolation + weight = index - lower + return sorted_data[lower] * (1 - weight) + sorted_data[upper] * weight + + return { + 'p50_ms': get_percentile(50), + 'p95_ms': get_percentile(95), + 'p99_ms': get_percentile(99), + } + + +# ============================================================================ +# Trace Loader +# ============================================================================ + +class TraceLoader: + """Load KVCache trace""" + + def __init__(self, trace_path: str): + """Initialize trace loader + + Args: + trace_path: Trace file path + """ + self.trace_path = trace_path + self.requests = [] + self._load_trace() + + def _load_trace(self): + """Load trace file with error handling""" + line_num = 0 + try: + with open(self.trace_path, 'r') as f: + for line in f: + line_num += 1 + line = line.strip() + if not line: + continue + try: + req = json.loads(line) + # Validate required fields + if not all(k in req for k in ['timestamp', 'hash_ids', 'input_length', 'output_length']): + print(f"Warning: Line {line_num} missing required fields, skipping") + continue + if not isinstance(req['hash_ids'], list): + print(f"Warning: Line {line_num} has invalid hash_ids (not a list), skipping") + continue + self.requests.append(KVCacheRequest( + timestamp=float(req['timestamp']), + hash_ids=req['hash_ids'], + input_length=int(req['input_length']), + output_length=int(req['output_length']) + )) + except (json.JSONDecodeError, ValueError, KeyError) as e: + print(f"Warning: Line {line_num} has invalid format: {e}, skipping") + continue + except FileNotFoundError: + raise FileNotFoundError(f"Trace file not found: {self.trace_path}") + except OSError as e: + raise OSError(f"Error reading trace file {self.trace_path}: {e}") + + def get_requests(self) -> List[KVCacheRequest]: + """Get request list + + Returns: + List[KVCacheRequest]: Request list + """ + return self.requests + + +# ============================================================================ +# Benchmark Runner +# ============================================================================ + +def run_benchmark(trace_path: str, storage_dir: str, bytes_per_token: int = DEFAULT_BYTES_PER_TOKEN, + max_requests: Optional[int] = None, max_blocks: int = 100000, + replay_timestamps: bool = False, time_scale: float = 1.0, + block_size_tokens: int = 512, + fsync_mode: str = 'batch', fsync_batch_size: int = 100) -> Dict: + """Run benchmark + + Args: + trace_path: Trace file path + storage_dir: Storage directory + bytes_per_token: Bytes per token + max_requests: Maximum number of requests (None = all) + max_blocks: Maximum number of blocks + replay_timestamps: Whether to replay timestamps from trace (simulate realistic timing) + time_scale: Time scaling factor (1.0=real-time, 0.1=10x speed, 10.0=0.1x speed) + block_size_tokens: Number of tokens per block + fsync_mode: When to fsync ('batch', 'always', 'end', 'none') + fsync_batch_size: Number of writes between fsync in batch mode + + Returns: + Dict: Benchmark results + """ + block_size_bytes = block_size_tokens * bytes_per_token + + print(f"\n{'='*80}") + print(f"Running: {Path(trace_path).name}") + print(f"Architecture: Offset Allocator (Mooncake style)") + print(f"Block size: {block_size_tokens} tokens/block ({block_size_bytes:,} bytes)") + print(f"Storage: Single large file with offset-based block management") + print(f"Bytes per token: {bytes_per_token}") + print(f"Max blocks: {max_blocks}") + print(f"Fsync mode: {fsync_mode}" + (f" (batch_size={fsync_batch_size})" if fsync_mode == 'batch' else '')) + print(f"Timestamp replay: {'Enabled' if replay_timestamps else 'Disabled'}") + if replay_timestamps: + scale_desc = 'real-time' if time_scale == 1.0 else f'{1/time_scale:.1f}x speed' if time_scale < 1.0 else f'{time_scale}x slower' + print(f"Time scale: {time_scale}x ({scale_desc})") + print(f"{'='*80}") + + # Load trace + loader = TraceLoader(trace_path) + requests = loader.get_requests() + + if max_requests: + requests = requests[:max_requests] + + print(f"Loaded {len(requests)} requests") + + # Show timestamp range + if replay_timestamps and requests: + timestamps = [req.timestamp for req in requests] + time_span_ms = max(timestamps) - min(timestamps) + print(f"Timestamp range: {min(timestamps):.1f} - {max(timestamps):.1f} ms (span: {time_span_ms:.1f} ms)") + + # Create benchmark instance with context manager for cleanup + with StorageBenchmark( + storage_dir, bytes_per_token, max_blocks, + block_size_tokens, fsync_mode, fsync_batch_size + ) as benchmark: + + # Run benchmark + start_time = time.perf_counter() + total_io_time = 0.0 # Actual I/O time (excluding sleep) + last_timestamp = None + base_time = time.time() # Use wall time for replay synchronization + + for i, req in enumerate(requests): + # Replay by timestamps + sleep_time = 0.0 + if replay_timestamps and last_timestamp is not None: + # Calculate time interval from previous request + delta_ms = req.timestamp - last_timestamp + sleep_time = delta_ms / 1000.0 / time_scale # Apply time scaling + + if sleep_time > 0: + time.sleep(sleep_time) + + # Process request (measure I/O time) + req_start = time.perf_counter() + benchmark.process_request(req) + req_io_time = time.perf_counter() - req_start + total_io_time += req_io_time + + # Record current request timestamp + last_timestamp = req.timestamp + + # Progress output + if (i + 1) % 100 == 0: + if replay_timestamps: + elapsed_wall_time = time.time() - base_time + simulated_time = (req.timestamp - requests[0].timestamp) / 1000.0 / time_scale + print(f" Processed {i + 1}/{len(requests)}... (wall: {elapsed_wall_time:.1f}s, simulated: {simulated_time:.1f}s, io: {total_io_time:.1f}s)") + else: + print(f" Processed {i + 1}/{len(requests)}...") + + elapsed = time.perf_counter() - start_time + + # Perform final sync to include it in stats + benchmark.storage._finalize_sync() + + # Get statistics (context manager will handle cleanup) + stats = benchmark.get_stats() + + # Calculate actual I/O time (excluding sleep) + io_time = total_io_time if replay_timestamps else elapsed + + return { + 'trace_file': Path(trace_path).name, + 'total_requests': len(requests), + 'simulation_time_s': elapsed, + 'io_time_s': io_time, # Actual I/O time + 'wall_time_s': elapsed, # Wall time (including sleep) + 'requests_per_second': len(requests) / io_time if io_time > 0 else 0, # Based on I/O time + 'timestamp_replay_enabled': replay_timestamps, + 'time_scale': time_scale, + 'bytes_per_token': bytes_per_token, + 'block_size_tokens': block_size_tokens, + 'fsync_mode': fsync_mode, + **stats, + } + + +# ============================================================================ +# Result Output +# ============================================================================ + +def print_results(results: List[Dict]): + """Print benchmark results + + Args: + results: List of benchmark results + """ + for i, r in enumerate(results, 1): + print(f"\n{'='*80}") + print(f" [{i}/{len(results)}] {r['trace_file']}") + print(f"{'='*80}") + + print(f"\n[Performance Overview]") + print(f" Total Requests: {r['total_requests']:,}") + print(f" Queries Per Second (QPS): {r['requests_per_second']:.2f}") + print(f" Cache Hit Rate: {r['block_hit_rate']:.2%}") + print(f" Write Ratio: {r['write_ratio']:.2%}") + print(f" Total Blocks: {r['total_blocks']:,}") + print(f" Read Blocks: {r['read_blocks']:,}") + print(f" Write Blocks: {r['write_blocks']:,}") + print(f" Prefix Hits: {r['prefix_hit_blocks']:,}") + + print(f"\n[Latency Analysis]") + req_lat = r['latency'] + print(f" Request Latency (End-to-End): Avg={req_lat['avg_ms']:.2f}ms, P50={req_lat['p50_ms']:.2f}ms, P95={req_lat['p95_ms']:.2f}ms, P99={req_lat['p99_ms']:.2f}ms") + read_lat = r['storage']['read'] + write_lat = r['storage']['write'] + print(f" Single I/O Operation (Per Block):") + print(f" Read: Avg={read_lat.get('avg_ms', 0):.3f}ms, P50={read_lat.get('p50_ms', 0):.3f}ms, P95={read_lat.get('p95_ms', 0):.3f}ms, P99={read_lat.get('p99_ms', 0):.3f}ms") + print(f" Write: Avg={write_lat.get('avg_ms', 0):.3f}ms, P50={write_lat.get('p50_ms', 0):.3f}ms, P95={write_lat.get('p95_ms', 0):.3f}ms, P99={write_lat.get('p99_ms', 0):.3f}ms") + + print(f"\n[I/O & Bandwidth]") + print(f" Total Read I/O: {r['storage']['read']['mb']:>10.1f} MB ({r['storage']['read']['count']:,} ops)") + print(f" Total Write I/O: {r['storage']['write']['mb']:>10.1f} MB ({r['storage']['write']['count']:,} ops)") + io_time = r['io_time_s'] + bandwidth = (r['storage']['read']['mb'] + r['storage']['write']['mb']) / io_time + print(f" Effective Bandwidth: {bandwidth:>10.1f} MB/s") + + print(f"\n[Storage Details]") + print(f" Blocks in Use: {r['storage']['total_blocks']:>10,}") + print(f" Free Blocks: {r['storage']['free_blocks']:>10,}") + print(f" Tokens per Block: {r['tokens_per_block']:>10,}") + print(f" Block Size: {r['tokens_per_block'] * r.get('bytes_per_token', 2048) / 1024 / 1024:>10.2f} MB") + if 'sync_count' in r['storage']: + print(f" Fsync Operations: {r['storage']['sync_count']:>10,}") + + print(f"\n[Execution Time]") + if r.get('timestamp_replay_enabled'): + print(f" Wall Time (Total): {r['wall_time_s']:>10.2f} s") + print(f" I/O Time (Actual): {r['io_time_s']:>10.2f} s") + print(f" Sleep Time (Replay): {r['wall_time_s'] - r['io_time_s']:>10.2f} s") + else: + print(f" Total Execution Time: {r['wall_time_s']:>10.2f} s") + + print(f"\n{'='*80}\n") + + +# ============================================================================ +# Main Program +# ============================================================================ + +def main(): + """Main entry point""" + parser = argparse.ArgumentParser( + description='Mooncake KVCache Storage Benchmark', + formatter_class=argparse.RawDescriptionHelpFormatter, + epilog=""" +Examples: + # Quick test (100 requests) + python storage_benchmark.py --scenario=toolagent --max-requests=100 + + # Test with large model preset (Llama-3.1-405B) + python storage_benchmark.py --scenario=toolagent --model=llama-3.1-405b --max-requests=100 + + # Test with Deepseek V3 (extra large model) + python storage_benchmark.py --scenario=toolagent --model=deepseek-v3 --max-requests=100 + + # Realistic replay (with timestamps, 10x speed) + python storage_benchmark.py --scenario=toolagent --max-requests=1000 \\ + --replay-timestamps --time-scale=0.1 + + # All scenarios with custom bytes_per_token + python storage_benchmark.py --scenario=all --bytes-per-token=512 + + # Test with different block sizes and fsync modes + python storage_benchmark.py --scenario=toolagent --block-size-tokens=256 --fsync-mode=always + + # Test with custom fsync batch size + python storage_benchmark.py --scenario=toolagent --fsync-mode=batch --fsync-batch-size=50 + +Performance Tuning: + --fsync-mode=batch (default): Balance between performance and safety + --fsync-mode=always: Safest but slowest, measures full persistence cost + --fsync-mode=end: Fastest, only measures write I/O (not persistence) + --fsync-mode=none: Testing only, no durability guarantees + +Available model presets: + llama-3.1-405b, qwen3-32b, deepseek-v3, glm-4.6, default + +For more information: tools/STORAGE_BENCHMARK_README.md + """ + ) + + parser.add_argument('--trace-dir', type=str, default='../../FAST25-release/traces', + help='Trace files directory') + parser.add_argument('--scenario', type=str, choices=['conversation', 'synthetic', 'toolagent', 'all'], + default='toolagent', help='Test scenario') + parser.add_argument('--storage-dir', type=str, default='/tmp/mooncake_bench', + help='Storage directory') + parser.add_argument('--model', type=str, choices=list(MODEL_BYTES_PER_TOKEN.keys()), + default='default', + help=f'Model preset (overrides --bytes-per-token). Available: {", ".join(MODEL_BYTES_PER_TOKEN.keys())}') + parser.add_argument('--bytes-per-token', type=int, default=DEFAULT_BYTES_PER_TOKEN, + help='Bytes per token (default %d, overridden by --model if specified)' % DEFAULT_BYTES_PER_TOKEN) + parser.add_argument('--max-requests', type=int, default=None, + help='Maximum number of requests (default: unlimited)') + parser.add_argument('--max-blocks', type=int, default=100000, + help='Maximum number of blocks in storage file (determines file size)') + parser.add_argument('--replay-timestamps', action='store_true', + help='Enable timestamp replay (simulate realistic request timing)') + parser.add_argument('--time-scale', type=float, default=1.0, + help='Time scaling factor (1.0=real-time, 0.1=10x speed, 10.0=0.1x speed)') + parser.add_argument('--block-size-tokens', type=int, default=512, + help='Number of tokens per block (default: 512)') + parser.add_argument('--fsync-mode', type=str, choices=['batch', 'always', 'end', 'none'], + default='batch', + help='When to fsync: batch=every N writes (default), always=after each write, end=only at close, none=never') + parser.add_argument('--fsync-batch-size', type=int, default=100, + help='Number of writes between fsync in batch mode (default: 100)') + + args = parser.parse_args() + + # Print benchmark header + print(f"\n{'='*80}") + print(f"{'Mooncake KVCache Storage Benchmark':^80}") + print(f"{'='*80}") + + # Determine bytes_per_token (model preset takes precedence) + bytes_per_token = MODEL_BYTES_PER_TOKEN.get(args.model, args.bytes_per_token) + if args.model != 'default': + print(f"Using model preset: {args.model} ({bytes_per_token} bytes/token, ~{bytes_per_token/1024:.1f} KB/token)") + else: + print(f"Using custom bytes_per_token: {bytes_per_token}") + + # Determine test scenarios + scenarios = ['conversation', 'synthetic', 'toolagent'] if args.scenario == 'all' else [args.scenario] + trace_files = { + 'conversation': 'conversation_trace.jsonl', + 'synthetic': 'synthetic_trace.jsonl', + 'toolagent': 'toolagent_trace.jsonl' + } + + # Run benchmarks + results = [] + + for scenario in scenarios: + trace_path = Path(args.trace_dir) / trace_files[scenario] + if trace_path.exists(): + result = run_benchmark( + str(trace_path), + str(Path(args.storage_dir) / scenario), + bytes_per_token, + args.max_requests, + args.max_blocks, + args.replay_timestamps, + args.time_scale, + args.block_size_tokens, + args.fsync_mode, + args.fsync_batch_size + ) + results.append(result) + else: + print(f"Warning: Trace file not found: {trace_path}") + + # Print results + if results: + print_results(results) + + +if __name__ == '__main__': + main() diff --git a/docs/source/performance/storage-benchmark.md b/docs/source/performance/storage-benchmark.md new file mode 100644 index 00000000..2ebabcdc --- /dev/null +++ b/docs/source/performance/storage-benchmark.md @@ -0,0 +1,163 @@ +# Mooncake KVCache Storage Benchmark + +High-performance KVCache storage benchmark tool based on Mooncake Store architecture. + +## Overview + +Evaluates I/O performance of KVCache storage systems using: +- Single large file (100GB) with offset-based block management +- Prefix caching simulation with hash-based block lookup +- Timestamp-based request replay for realistic testing +- Comprehensive metrics: latency, bandwidth, hit rates + +## Test Flow + +1. **Load Traces**: Read request sequences from JSONL files (`FAST25-release/traces`) +2. **Process Requests**: For each request, check hash_id prefix cache hits/misses +3. **Perform I/O**: Read cached blocks from disk, write new blocks to storage +4. **Collect Metrics**: Track latency, bandwidth, and cache hit rates + +## Quick Start + +```bash +# Quick test (100 requests, no timestamp replay) +python storage_benchmark.py --scenario=toolagent --max-requests=100 + +# Test with large model preset (Llama-3.1-405B) +python storage_benchmark.py --scenario=toolagent --model=llama-3.1-405b --max-requests=100 + +# Test with Deepseek V3 (extra large model) +python storage_benchmark.py --scenario=toolagent --model=deepseek-v3 --max-requests=100 + +# Realistic replay (with timestamps, 10x speed) +python storage_benchmark.py --scenario=toolagent --max-requests=1000 \ + --replay-timestamps --time-scale=0.1 + +# Test all scenarios with replay +python storage_benchmark.py --scenario=all --time-scale=1.0 +``` + +## Command-Line Options + +| Option | Description | Default | +|--------|-------------|---------| +| `--trace-dir` | Trace files directory | `../FAST25-release/traces` | +| `--scenario` | Test scenario: `conversation`, `synthetic`, `toolagent`, `all` | `toolagent` | +| `--storage-dir` | Storage directory | `/tmp/mooncake_bench` | +| `--model` | Model preset (overrides `--bytes-per-token`) | `default` | +| `--bytes-per-token` | Bytes per token (2048 for 7B FP16) | `2048` | +| `--max-requests` | Maximum requests per scenario (unlimited if not specified) | `None` | +| `--max-blocks` | Maximum number of blocks | `100000` | +| `--replay-timestamps` | Enable timestamp replay | `False` | +| `--time-scale` | Time scaling factor (1.0 = real-time, 0.1 = 10x faster) | `1.0` | + +## Model Presets + +The tool includes presets for popular LLM models with accurate KVCache sizes based on the [LMCache KVCache Calculator](https://lmcache.ai/kv_cache_calculator.html). + +| Model | Bytes/Token | Size | Notes | +|-------|-------------|------|-------| +| **Small Models (7B-13B)** | +| `llama-3-8b` | 128 | 128 B/token | GQA optimized | +| `mistral-7b` | 128 | 128 B/token | GQA optimized | +| `qwen-14b` | 40 | 40 B/token | GQA optimized | +| `gemma-7b` | 224 | 224 B/token | | +| `llama-2-7b` | 512 | 512 B/token | | +| `llama-2-13b` | 800 | 800 B/token | | +| **Large Models (70B-405B)** | +| `llama-2-70b` | 320 | 320 B/token | GQA optimized | +| `llama-3-70b` | 320 | 320 B/token | GQA optimized | +| `mixtral-8x7b` | 128 | 128 B/token | GQA optimized | +| `mixtral-8x22b` | 224 | 224 B/token | GQA optimized | +| `qwen-72b` | 320 | 320 B/token | GQA optimized | +| `qwen-110b` | 320 | 320 B/token | GQA optimized | +| `llama-3.1-405b` | 516018 | ~504 KB/token | Very large KVCache | +| **Extra Large Models** | +| `glm-4.6` | 156991 | ~153 KB/token | | +| `deepseek-v3` | 1749384 | ~1.67 MB/token | Largest KVCache | +| **Default** | +| `default` | 2048 | 2 KB/token | Legacy 7B FP16 | + +**Usage**: `--model=llama-3.1-405b` (overrides `--bytes-per-token`) + +## Test Scenarios + +- **`conversation`**: Write-intensive workload (dialogue patterns) +- **`synthetic`**: Read-intensive workload (cached patterns) +- **`toolagent`**: Balanced read/write mix (tool use patterns) + +## Output Example + +``` +================================================================================ +Mooncake KVCache Storage Benchmark +================================================================================ +Using model preset: llama-3.1-405b (516018 bytes/token, ~504.0 KB/token) + +[1/1] toolagent_trace.jsonl +================================================================================ + +[Performance Overview] + Total Requests: 100 + Queries Per Second (QPS): 14.45 + Cache Hit Rate: 24.27% + Write Ratio: 75.73% + Total Blocks: 1,949 + Read Blocks: 473 + Write Blocks: 1,476 + Prefix Hits: 376 + +[Latency Analysis] + Request Latency (End-to-End): Avg=69.18ms, P50=15.49ms, P95=239.99ms, P99=310.58ms + Single I/O Operation (Per Block): + Read: Avg=14.572ms, P50=0.280ms, P95=0.280ms, P99=0.280ms + Write: Avg=5.120ms, P50=5.120ms, P95=5.120ms, P99=5.120ms + +[I/O & Bandwidth] + Total Read I/O: 473.0 MB (473 ops) + Total Write I/O: 1476.0 MB (1,476 ops) + Effective Bandwidth: 280.8 MB/s + +[Storage Details] + Blocks in Use: 1,476 + Free Blocks: 0 + Tokens per Block: 512 + Block Size: 1.00 MB + +[Execution Time] + Total Execution Time: 8.42 s + +================================================================================ +``` + +## Metrics + +| Metric | Description | +|--------|-------------| +| **QPS** | Queries per second (based on I/O time, excluding sleep) | +| **Request Latency** | End-to-end latency for entire request (all I/O operations) | +| **Single I/O Latency** | Latency for individual block read/write operations (512 tokens) | +| **P50/P95/P99** | Latency percentiles (milliseconds) using linear interpolation | +| **Hit Rate** | Cache hit ratio for blocks | +| **Write Ratio** | Percentage of blocks that needed to be written | +| **Bandwidth** | Effective throughput based on I/O time only | +| **Prefix Hits** | Number of blocks served from prefix cache | + +**Note**: Request Latency measures the total time to process all blocks in a request, while Single I/O Latency measures the time for one block operation (512 tokens). + +## Trace Data Format + +```json +{ + "timestamp": 1234.567, + "hash_ids": [1, 2, 4, 7], + "input_length": 2048, + "output_length": 512 +} +``` + +Each `hash_id` corresponds to a 512-token block. The tool simulates prefix caching by checking if blocks are already in storage before writing. + +## Requirements + +- Python 3.10+