4.1-GPU显存分配(2)


return

# 按地址排序已分配块 sorted_blocks = sorted(self.allocated_blocks.values(), key=lambda x: x['start']) # 重新分配内存 new_start = 0 for block in sorted_blocks: old_start = block['start'] size = block['size'] # 移动块 if old_start != new_start: self._move_block(old_start, new_start, size) # 更新块信息 block['start'] = new_start new_start += size # 更新空闲块 remaining_memory = self.total_memory - new_start self.free_blocks = [(new_start, remaining_memory)] def _move_block(self, old_start, new_start, size): """移动内存块""" # 在实际应用中,这里需要实现具体的内存移动逻辑 # 这通常涉及到GPU内存拷贝操作 # 模拟内存移动 print(f"Moving block from {old_start} to {new_start}, size {size}") def analyze_compaction_benefit(self): """分析压缩收益""" if not self.allocated_blocks: return 0 # 计算移动开销 total_allocated = sum(block['size'] for block in self.allocated_blocks.values()) move_cost = total_allocated # 假设移动成本等于数据大小 # 计算碎片减少 initial_fragmentation = self._calculate_fragmentation() self.compact_memory() final_fragmentation = self._calculate_fragmentation() fragmentation_reduction = initial_fragmentation - final_fragmentation # 计算净收益 net_benefit = fragmentation_reduction - move_cost / self.total_memory return net_benefit def _calculate_fragmentation(self): """计算碎片程度""" if not self.allocated_blocks: return 0 total_allocated = sum(block['size'] for block in self.allocated_blocks.values()) total_free = self.total_memory - total_allocated if total_free == 0: return 0 # 简化计算:碎片程度 = (空闲块数量 - 1) / 总块数 free_block_count = len(self.free_blocks) total_block_count = free_block_count + len(self.allocated_blocks) fragmentation = (free_block_count - 1) / total_block_count if total_block_count > 0 else 0 return fragmentation
### 3.2 内存重分配策略 ```python class MemoryReallocator: def __init__(self, total_memory): self.total_memory = total_memory self.allocator = FirstFitAllocator(total_memory) self.compaction_threshold = 0.3 # 碎片超过30%时触发压缩 def reallocate_with_compaction(self, new_size): """带压缩的重分配""" # 检查当前碎片程度 if self._should_compact(): print("执行内存压缩...") self._compact_memory() # 尝试分配 result = self.allocator.allocate(new_size) if result is None: # 分配失败,尝试强制压缩后重试 self._force_compact() result = self.allocator.allocate(new_size) return result def _should_compact(self): """判断是否需要压缩""" # 简化判断:如果碎片块过多,需要压缩 allocated_count = len(self.allocator.allocated_blocks) free_count = len(self.allocator.free_blocks) total_blocks = allocated_count + free_count # 如果空闲块数量大于已分配块数量的1.5倍,认为碎片严重 return free_count > allocated_count * 1.5 def _compact_memory(self): """内存压缩""" compactor = MemoryCompactor(self.total_memory) compactor.allocated_blocks = self.allocator.allocated_blocks.copy() compactor.free_blocks = self.allocator.free_blocks.copy() compactor.compact_memory() # 更新分配器状态 self.allocator.allocated_blocks = compactor.allocated_blocks self.allocator.free_blocks = compactor.free_blocks def _force_compact(self): """强制压缩""" print("执行强制压缩...") self._compact_memory() # 清理一些临时块 self._cleanup_temporary_blocks() def _cleanup_temporary_blocks(self): """清理临时块""" # 在实际应用中,这里需要识别并释放临时块 print("清理临时块...")

3.3 增量式碎片整理

class IncrementalMemoryCompactor: def __init__(self, total_memory, max_compaction_time=100): self.total_memory = total_memory self.max_compaction_time = max_compaction_time # 毫秒 self.allocator = FirstFitAllocator(total_memory) self.compaction_queue = [] self.compaction_in_progress = False def schedule_compaction(self, priority='low'): """计划碎片整理""" self.compaction_queue.append({ 'priority': priority, 'timestamp': time.time(), 'completed': False }) def perform_incremental_compaction(self): """执行增量碎片整理""" if not self.compaction_queue or self.compaction_in_progress: return # 选择最高优先级的任务 self.compaction_queue.sort(key=lambda x: (x['priority'] != 'high', x['timestamp'])) task = self.compaction_queue[0] if task['completed']: self.compaction_queue.pop(0) return # 开始碎片整理 self.compaction_in_progress = True start_time = time.time() try: # 执行部分碎片整理 blocks_processed = self._partial_compaction(start_time) # 检查是否完成 if blocks_processed >= self._get_total_blocks(): task['completed'] = True print(f"碎片整理完成,处理了 {blocks_processed} 个块") else: print(f"碎片整理进行中,处理了 {blocks_processed} 个块") except Exception as e: print(f"碎片整理出错: {e}") finally: self.compaction_in_progress = False def _partial_compaction(self, start_time): """部分碎片整理""" processed_blocks = 0 # 按时间顺序处理已分配块 allocated_blocks = list(self.allocator.allocated_blocks.values()) for i, block in enumerate(allocated_blocks): # 检查时间限制 if (time.time() - start_time) * 1000 > self.max_compaction_time: break # 处理当前块 if self._should_move_block(block): self._move_block(block) processed_blocks += 1 return processed_blocks def _should_move_block(self, block): """判断是否需要移动块""" # 简化逻辑:如果块位于内存后半部分,且较大,则考虑移动 return block['start'] > self.total_memory * 0.6 and block['size'] > 1024*1024 def _move_block(self, block): """移动块""" # 在实际应用中,这里需要实现具体的内存移动逻辑 print(f"移动块从 {block['start']} 到新位置,大小 {block['size']}") def _get_total_blocks(self): """获取总块数""" return len(self.allocator.allocated_blocks) + len(self.allocator.free_blocks)

4. PyTorch中的显存管理

4.1 PyTorch显存管理机制

class PyTorchMemoryManager: def __init__(self, device='cuda'): self.device = device self.cuda_available = torch.cuda.is_available() self.current_device = torch.cuda.current_device() if self.cuda_available else 0 def get_memory_info(self): """获取PyTorch显存信息""" if not self.cuda_available: return None # 获取各种显存使用情况 allocated = torch.cuda.memory_allocated(self.current_device) reserved = torch.cuda.memory_reserved(self.current_device) max_allocated = torch.cuda.max_memory_allocated(self.current_device) # 获取总显存 total_memory = torch.cuda.get_device_properties(self.current_device).total_memory return { 'allocated': allocated, 'reserved': reserved, 'max_allocated': max_allocated, 'total': total_memory, 'utilization': allocated / total_memory * 100, 'fragmentation': self._calculate_fragmentation(allocated, reserved) } def _calculate_fragmentation(self, allocated, reserved): """计算碎片程度""" if reserved == 0: return 0 # 碎片程度 = (已保留 - 已分配) / 已保留 fragmentation = (reserved - allocated) / reserved return fragmentation def clear_cache(self): """清除显存缓存""" if self.cuda_available: torch.cuda.empty_cache() print("显存缓存已清除") def memory_monitor(self, interval=1): """内存监控""" print("开始显存监控...") try: while True: memory_info = self.get_memory_info() if memory_info: print(f"已分配: {memory_info['allocated']/1024**3:.2f}GB, " f"已保留: {memory_info['reserved']/1024**3:.2f}GB, " f"利用率: {memory_info['utilization']:.1f}%") time.sleep(interval) except KeyboardInterrupt: print("监控停止") def optimize_memory_usage(self): """优化内存使用""" if not self.cuda_available: return # 启用内存重用 torch.cuda.set_per_process_memory_fraction(0.9) # 使用90%显存 # 启用内存池 if torch.cuda.is_available(): torch.cuda.enable_mem_pool() # 清理不必要的缓存 self.clear_cache() print("显存优化完成")

4.2 自定义内存分配器

class CustomMemoryAllocator: def __init__(self, device='cuda', strategy='adaptive'): self.device = device self.strategy = strategy self.allocator = self._create_allocator() # 内存池 self.memory_pool = {} self.allocation_stats = { 'total_allocated': 0, 'total_freed': 0, 'peak_memory': 0, 'fragmentation_count': 0 } def _create_allocator(self): """创建分配器""" if self.strategy == 'first_fit': return FirstFitAllocator(16*1024**3) elif self.strategy == 'best_fit': return BestFitAllocator(16*1024**3) elif self.strategy == 'buddy': return BuddySystemAllocator(16*1024**3) else: return TieredMemoryAllocator(self.device) def allocate_tensor(self, size, dtype=torch.float32): """分配张量""" # 检查内存池 pool_key = (size, dtype) if pool_key in self.memory_pool: tensor = self.memory_pool[pool_key] del self.memory_pool[pool_key] return tensor # 使用分配器分配 start, allocated_size, tier = self.allocator.allocate(size) if start is None: raise MemoryError("内存分配失败") # 创建张量 tensor = torch.empty(size, dtype=dtype, device=self.device) # 记录分配 self.allocation_stats['total_allocated'] += allocated_size self.allocation_stats['peak_memory'] = max( self.allocation_stats['peak_memory'], self.allocation_stats['total_allocated'] - self.allocation_stats['total_freed'] ) return tensor def deallocate_tensor(self, tensor): """释放张量""" size = tensor.numel() * tensor.element_size() # 尝试放入内存池 pool_key = (tensor.size(), tensor.dtype) if len(self.memory_pool) < 100: # 限制内存池大小 self.memory_pool[pool_key] = tensor return # 直接释放 del tensor self.allocation_stats['total_freed'] += size def get_memory_stats(self): """获取内存统计""" return { **self.allocation_stats, 'current_memory': self.allocation_stats['total_allocated'] - self.allocation_stats['total_freed'], 'memory_pool_size': len(self.memory_pool) } def optimize_memory_pool(self): """优化内存池""" # 清理长时间未使用的内存池项 current_time = time.time() expired_items = [] for key, tensor in self.memory_pool.items(): # 检查内存是否仍然有效 try: if not tensor.is_contiguous(): expired_items.append(key) except: expired_items.append(key) # 清理过期项 for key in expired_items: del self.memory_pool[key] print(f"清理了 {len(expired_items)} 个内存池项")

4.3 训练中的内存优化

class TrainingMemoryOptimizer: def __init__(self, model, device='cuda'): self.model = model.to(device) self.device = device self.optimizer = None self.grad_scaler = None self.memory_manager = PyTorchMemoryManager(device) def setup_training(self, learning_rate=1e-4): """设置训练""" # 创建优化器 self.optimizer = torch.optim.AdamW(self.model.parameters(), lr=learning_rate) # 创建梯度缩放器(用于混合精度) if torch.cuda.is_available(): self.grad_scaler = torch.cuda.amp.GradScaler() print("训练设置完成") def train_with_memory_optimization(self, dataloader, epochs=10): """带内存优化的训练""" # 启用内存优化 self.enable_memory_optimization() for epoch in range(epochs): self.train_epoch(dataloader) # 定期检查内存使用 if epoch % 5 == 0: self.check_memory_usage() def enable_memory_optimization(self): """启用内存优化""" # 启用梯度检查点 torch.utils.checkpoint.enable() # 启用内存池 torch.cuda.enable_mem_pool() # 优化批处理大小 self.optimize_batch_size() print("内存优化已启用") def train_epoch(self, dataloader): """训练一个epoch""" self.model.train() for batch_idx, (data, target) in enumerate(dataloader): data, target = data.to(self.device), target.to(self.device) # 清除梯度 self.optimizer.zero_grad() # 前向传播 with torch.cuda.amp.autocast(): output = self.model(data) loss = torch.nn.functional.cross_entropy(output, target) # 反向传播 if self.grad_scaler: self.grad_scaler.scale(loss).backward() self.grad_scaler.step(self.optimizer) self.grad_scaler.update() else: loss.backward() self.optimizer.step() # 定期清除缓存 if batch_idx % 50 == 0: torch.cuda.empty_cache() def optimize_batch_size(self): """优化批处理大小""" # 获取可用内存 memory_info = self.memory_manager.get_memory_info() available_memory = memory_info['total'] - memory_info['allocated'] # 简化的批大小计算 estimated_batch_size = int(available_memory / (100 * 1024 * 1024)) # 100MB per batch print(f"优化批大小为: {estimated_batch_size}") return estimated_batch_size def check_memory_usage(self): """检查内存使用""" memory_info = self.memory_manager.get_memory_info() print(f"Epoch {epoch}: " f"已分配: {memory_info['allocated']/1024**3:.2f}GB, " f"利用率: {memory_info['utilization']:.1f}%") # 如果内存使用过高,触发清理 if memory_info['utilization'] > 90: print("内存使用过高,触发清理...") torch.cuda.empty_cache()

5. 实际应用案例分析

5.1 大语言模型内存管理

class LLMMemoryManager: def __init__(self, model_name='gpt2', device='cuda'): self.model_name = model_name self.device = device self.model = None self.kv_cache = None self.memory_optimizer = None def load_model(self): """加载模型""" from transformers import AutoModelForCausalLM self.model = AutoModelForCausalLM.from_pretrained(self.model_name) self.model = self.model.to(self.device) # 初始化KV缓存 self.initialize_kv_cache() # 创建内存优化器 self.memory_optimizer = TrainingMemoryOptimizer(self.model, self.device) print("模型加载完成") def initialize_kv_cache(self): """初始化KV缓存""" self.kv_cache = {} cache_size = 1024 * 1024 * 1024 # 1GB self.kv_cache_memory = TieredMemoryAllocator(self.device) def generate_text(self, prompt, max_length=100): """生成文本""" # 预处理 inputs = self._preprocess_input(prompt) # 生成文本 outputs = self.model.generate( inputs, max_length=max_length, num_beams=1, do_sample=True, temperature=0.7, pad_token_id=self.model.config.eos_token_id ) return self._post_process_output(outputs) def _preprocess_input(self, prompt): """预处理输入""" # 在实际应用中,这里需要实现tokenization return torch.tensor([[1, 2, 3]]) # 简化的预处理 def _post_process_output(self, outputs): """后处理输出""" # 在实际应用中,这里需要实现反tokenization return "Generated text" def optimize_for_inference(self): """优化推理""" # 量化模型 self.quantize_model() # 优化KV缓存 self.optimize_kv_cache() # 启用内存优化 self.memory_optimizer.enable_memory_optimization() print("推理优化完成") def quantize_model(self): """量化模型""" # 简化的量化实现 for param in self.model.parameters(): if param.dtype == torch.float32: param.data = param.data.to(torch.float16) print("模型量化完成") def optimize_kv_cache(self): """优化KV缓存""" # 使用分级内存分配器管理KV缓存 self.kv_cache_memory = TieredMemoryAllocator(self.device) print("KV缓存优化完成")

5.2 图像训练内存优化

class ImageTrainingMemoryOptimizer: def __init__(self, model

作者与出处
整理: 灏天文库整理
本站整理收录,版权归原作者/开源协议所有;欢迎通过原文链接访问源仓库。
发布者: 作者: 脉冲星学徒的小龙虾 转发
评论区 (0)
U