diff --git a/distorch_2.py b/distorch_2.py index e4788a9..18dc4f0 100644 --- a/distorch_2.py +++ b/distorch_2.py @@ -7,6 +7,7 @@ import sys import torch import logging import hashlib +import re logger = logging.getLogger("MultiGPU") import copy @@ -128,79 +129,247 @@ def register_patched_safetensor_modelpatcher(): logger.info("[MultiGPU_DisTorch2] Successfully patched ModelPatcher.partially_load") +def parse_memory_string(mem_str): + """Parses a memory string (e.g., '4.0g', '512M') and returns bytes.""" + mem_str = mem_str.strip().lower() + match = re.match(r'(\d+\.?\d*)\s*([gmkb]?)', mem_str) + if not match: + raise ValueError(f"Invalid memory string format: {mem_str}") + + val, unit = match.groups() + val = float(val) + + if unit == 'g': + return val * (1024**3) + elif unit == 'm': + return val * (1024**2) + elif unit == 'k': + return val * 1024 + else: # b or no unit + return val + def analyze_safetensor_loading(model_patcher, allocations_str): """ Analyze and distribute safetensor model blocks across devices """ - DEVICE_RATIOS_DISTORCH = {} - device_table = {} distorch_alloc = allocations_str - virtual_vram_gb = 0.0 - if '#' in allocations_str: - distorch_alloc, virtual_vram_str = allocations_str.split('#') - if not distorch_alloc: - distorch_alloc = calculate_safetensor_vvram_allocation(model_patcher, virtual_vram_str) + distorch_alloc, _ = allocations_str.split('#', 1) eq_line = "=" * 50 dash_line = "-" * 50 fmt_assign = "{:<18}{:>7}{:>14}{:>10}" - for allocation in distorch_alloc.split(';'): - if ',' not in allocation: - continue - dev_name, fraction = allocation.split(',') - fraction = float(fraction) - total_mem_bytes = mm.get_total_memory(torch.device(dev_name)) - alloc_gb = (total_mem_bytes * fraction) / (1024**3) - DEVICE_RATIOS_DISTORCH[dev_name] = alloc_gb - device_table[dev_name] = { - "fraction": fraction, - "total_gb": total_mem_bytes / (1024**3), - "alloc_gb": alloc_gb - } + raw_block_list = model_patcher._load_list() + total_memory = sum(module_size for module_size, _, _, _ in raw_block_list) + + mode = "fraction" + if any(c in distorch_alloc.lower() for c in ['g', 'm', 'k', 'b']): + mode = "byte" + elif "%" in distorch_alloc: + mode = "ratio" + parsed_allocations = {} + wildcard_device = "cpu" # Default + + # Parse allocation string + raw_parsed = {} + user_requested_values = {} + for allocation in distorch_alloc.split(';'): + if ',' not in allocation: continue + dev_name, val_str = allocation.split(',', 1) + if '*' in dev_name: + dev_name = dev_name.replace('*','').strip() + wildcard_device = dev_name + + try: + if mode == "ratio": + value = float(val_str.replace('%','').strip()) + raw_parsed[dev_name] = value + user_requested_values[dev_name] = f"{value:.1f}%" + elif mode == "byte": + value_bytes = parse_memory_string(val_str) + raw_parsed[dev_name] = value_bytes + user_requested_values[dev_name] = f"{value_bytes / (1024**3):.2f}g" + else: # fraction + fraction = float(val_str) + total_dev_mem = mm.get_total_memory(torch.device(dev_name)) + parsed_allocations[dev_name] = total_dev_mem * fraction + user_requested_values[dev_name] = f"{int(fraction * 100)}%" + except ValueError as e: + logger.error(f"[MultiGPU_DisTorch2] Could not parse allocation '{allocation}': {e}") + return + + # Normalize and finalize allocations + if mode in ["ratio", "byte"]: + total_requested = sum(raw_parsed.values()) + if mode == "ratio": + for dev, val in raw_parsed.items(): + parsed_allocations[dev] = (val / total_requested) * total_memory + elif mode == "byte": + if total_requested > total_memory: + logger.info(f"[MultiGPU_DisTorch2] Over-allocation: Requested {total_requested/(1024**3):.2f}GB, but model is {total_memory/(1024**3):.2f}GB. Pro-rating allocations.") + for dev, val in raw_parsed.items(): + parsed_allocations[dev] = (val / total_requested) * total_memory + else: + parsed_allocations = raw_parsed + if wildcard_device not in parsed_allocations: + parsed_allocations[wildcard_device] = 0 + + remaining_mem = total_memory - total_requested + if remaining_mem > 0: + logger.info(f"[MultiGPU_DisTorch2] Under-allocation: {remaining_mem/(1024**2):.2f}MB of model unallocated. Assigning to wildcard device '{wildcard_device}'.") + parsed_allocations[wildcard_device] += remaining_mem + + if wildcard_device not in parsed_allocations: + parsed_allocations[wildcard_device] = 0 + + # Provide user feedback on allocation interpretation + if not distorch_alloc or distorch_alloc.isspace(): + logger.info("[MultiGPU_DisTorch2] Examples:") + logger.info(" Direct(byte) Mode - cuda:0,500mb;cuda:1,3.0g;cpu,5gb* -> '*' cpu = over/underflow device, put 0.50gb on cuda0, 3.00gb on cuda1, and 5.00gb (or the rest) on cpu") + logger.info(" Ratio(%) Mode - cuda:0,8%;cuda:1,8%;cpu,4% -> 8:8:4 ratio, put 40% on cuda0, 40% on cuda1, and 20% on cpu") + else: + if mode == "byte": + feedback_parts = [] + for dev in sorted_devices: + if dev in user_requested_values: + val_without_g = user_requested_values[dev].rstrip('g') + feedback_parts.append(f"{dev},{val_without_g}g") + else: + feedback_parts.append(f"{dev},0.00g") + + wildcard_indicator = "" + if wildcard_device: + wildcard_indicator = f"*{wildcard_device}" + + logger.info(f"[MultiGPU_DisTorch2] Interpreted Byte Allocation: {';'.join(feedback_parts)}{wildcard_indicator}") + + original_parts = [] + original_wildcard_device = None + for allocation in distorch_alloc.split(';'): + if ',' not in allocation: continue + dev_name, val_str = allocation.split(',', 1) + if '*' in dev_name: + dev_name = dev_name.replace('*','').strip() + original_wildcard_device = dev_name + original_parts.append((dev_name, val_str.strip())) + + if original_parts: + formatted_parts = [] + for dev_name, val_str in original_parts: + if 'mb' in val_str.lower(): + mb_val = float(val_str.lower().replace('mb', '')) + gb_val = mb_val / 1024 + formatted_parts.append(f"{gb_val:.2f}gb on {dev_name}") + elif 'gb' in val_str.lower() or 'g' in val_str.lower(): + val_num = float(''.join(filter(lambda x: x.isdigit() or x == '.', val_str))) + formatted_parts.append(f"{val_num:.2f}gb on {dev_name}") + else: + formatted_parts.append(f"{val_str} on {dev_name}") + + if formatted_parts: + if len(formatted_parts) == 1: + put_part = formatted_parts[0] + elif len(formatted_parts) == 2: + put_part = f"{formatted_parts[0]} and {formatted_parts[1]}" + else: + put_part = ", ".join(formatted_parts[:-1]) + f", and {formatted_parts[-1]}" + + wildcard_dev = original_wildcard_device if original_wildcard_device else "cpu" + logger.info(f"[MultiGPU_DisTorch2] Direct(byte) Mode - {distorch_alloc} -> '*' {wildcard_dev} = over/underflow device, put {put_part}") + + elif mode == "ratio": + total_requested_percent = sum(raw_parsed.values()) + if total_requested_percent > 0: + normalized_ratios = {} + for dev, percent in raw_parsed.items(): + normalized_ratios[dev] = (percent / total_requested_percent) * 100 + + ratio_parts = [] + ratio_values = [] + for allocation in distorch_alloc.split(';'): + if ',' not in allocation: continue + dev_name, val_str = allocation.split(',', 1) + val_str = val_str.strip() + if '%' in val_str: + ratio_val = val_str.replace('%', '').strip() + ratio_values.append(ratio_val) + ratio_parts.append(f"{dev_name}") + + if ratio_values and ratio_parts: + total_ratio = sum(float(val) for val in ratio_values) + normalized_pcts = [] + for val in ratio_values: + normalized_pct = (float(val) / total_ratio) * 100 + normalized_pcts.append(int(normalized_pct)) + + ratio_string = ":".join(ratio_values) + + put_parts = [] + for i, (dev_name, pct) in enumerate(zip(ratio_parts, normalized_pcts)): + put_parts.append(f"{pct}% on {dev_name}") + + if len(put_parts) == 1: + put_part = put_parts[0] + elif len(put_parts) == 2: + put_part = f"{put_parts[0]} and {put_parts[1]}" + else: + put_part = ", ".join(put_parts[:-1]) + f", and {put_parts[-1]}" + + logger.info(f"[MultiGPU_DisTorch2] Ratio(%) Mode - {distorch_alloc} -> {ratio_string} ratio, put {put_part}") + + # Log allocation table logger.info(eq_line) logger.info(" DisTorch2 Model Device Allocations") logger.info(eq_line) - logger.info(fmt_assign.format("Device", "Alloc %", "Total (GB)", " Alloc (GB)")) + + fmt_rosetta = "{:<10}{:>8}{:>8}{:>10}{:>10}" + logger.info(fmt_rosetta.format("Device", "VRAM GB", "Dev %", "Model GB", "Dist Ratio")) logger.info(dash_line) - sorted_devices = sorted(device_table.keys(), key=lambda d: (d == "cpu", d)) - - for dev in sorted_devices: - frac = device_table[dev]["fraction"] - tot_gb = device_table[dev]["total_gb"] - alloc_gb = device_table[dev]["alloc_gb"] - logger.info(fmt_assign.format(dev,f"{int(frac * 100)}%",f"{tot_gb:.2f}",f"{alloc_gb:.2f}")) + sorted_devices = sorted(parsed_allocations.keys(), key=lambda d: (d == "cpu", d)) + + dist_ratio_values = [] + if mode == "ratio": + total_requested_percent = sum(raw_parsed.values()) + for dev in sorted_devices: + if dev in raw_parsed: + normalized_pct = (raw_parsed[dev] / total_requested_percent) * 100 + dist_ratio_values.append(f"{int(normalized_pct)}%") + else: + dist_ratio_values.append("0%") + elif mode == "byte" or mode == "fraction": + for dev in sorted_devices: + model_percent = (parsed_allocations[dev] / total_memory) * 100 if total_memory > 0 else 0 + dist_ratio_values.append(f"{model_percent:.1f}%") + for i, dev in enumerate(sorted_devices): + alloc_gb = parsed_allocations[dev] / (1024**3) + total_dev_gb = mm.get_total_memory(torch.device(dev)) / (1024**3) + device_percent = (parsed_allocations[dev] / (total_dev_gb * 1024**3)) * 100 if total_dev_gb > 0 else 0 + + dist_ratio_str = dist_ratio_values[i] if i < len(dist_ratio_values) else "N/A" + + logger.info(fmt_rosetta.format(dev, f"{total_dev_gb:.2f}", f"{device_percent:.1f}%", f"{alloc_gb:.2f}", dist_ratio_str)) + logger.info(dash_line) + # Build block lists block_summary = {} - block_list = [] memory_by_type = defaultdict(int) - total_memory = 0 - - raw_block_list = model_patcher._load_list() - - # Calculate total memory from ComfyUI's list (first pass replacement) - total_memory = sum(module_size for module_size, _, _, _ in raw_block_list) - - # Set the minimum block size threshold (0.01% of total model memory) + MIN_BLOCK_THRESHOLD = total_memory * 0.0001 - logger.debug(f"[MultiGPU_DisTorch2] Total model memory: {total_memory} bytes") + logger.debug(f"[MultiGPU_DisTorch2] Total model memory: {total_memory / (1024**2):.2f} MB") logger.debug(f"[MultiGPU_DisTorch2] Tiny block threshold (0.01%): {MIN_BLOCK_THRESHOLD} bytes") - # Build all_blocks from ComfyUI's list (second pass replacement) all_blocks = [] for module_size, module_name, module_object, params in raw_block_list: block_type = type(module_object).__name__ - # Populate summary dictionaries block_summary[block_type] = block_summary.get(block_type, 0) + 1 memory_by_type[block_type] += module_size all_blocks.append((module_name, module_object, block_type, module_size)) - # Filter out tiny blocks from the distribution list block_list = [b for b in all_blocks if b[3] >= MIN_BLOCK_THRESHOLD] tiny_block_list = [b for b in all_blocks if b[3] < MIN_BLOCK_THRESHOLD] @@ -208,6 +377,7 @@ def analyze_safetensor_loading(model_patcher, allocations_str): logger.debug(f"[MultiGPU_DisTorch2] Distributable blocks: {len(block_list)}") logger.debug(f"[MultiGPU_DisTorch2] Tiny blocks (<0.01%): {len(tiny_block_list)}") + # Log layer distribution logger.info(" DisTorch2 Model Layer Distribution") logger.info(dash_line) fmt_layer = "{:<18}{:>7}{:>14}{:>10}" @@ -221,100 +391,124 @@ def analyze_safetensor_loading(model_patcher, allocations_str): logger.info(dash_line) - # Distribute blocks sequentially from the tail of the model - device_assignments = {device: [] for device in DEVICE_RATIOS_DISTORCH.keys()} + # Distribute blocks block_assignments = {} - - # Determine the primary compute device (first non-cpu device) - compute_device = "cuda:0" # Fallback + device_quotas = parsed_allocations.copy() + + compute_device = "cuda:0" for dev in sorted_devices: if dev != "cpu": compute_device = dev break - # Calculate total memory to be offloaded to donor devices - total_offload_gb = sum(DEVICE_RATIOS_DISTORCH.get(d, 0) for d in sorted_devices if d != compute_device) - total_offload_bytes = total_offload_gb * (1024**3) - - offloaded_bytes = 0 - - # Iterate from the TAIL of the model - for block_name, module, block_type, block_memory in reversed(block_list): - try: - # block_memory is already calculated - pass - except: - block_memory = 0 - if hasattr(module, 'weight') and module.weight is not None: - block_memory += module.weight.numel() * module.weight.element_size() - if hasattr(module, 'bias') and module.bias is not None: - block_memory += module.bias.numel() * module.bias.element_size() + if mode == 'gb': + devices_to_fill = [d for d in sorted_devices if d != wildcard_device] + else: + devices_to_fill = sorted(device_quotas.keys(), key=lambda d: (d == "cpu", d)) - # Assign to donor device (currently assumes one donor 'cpu') until target is met - if offloaded_bytes < total_offload_bytes: - # For now, simple offload to CPU, will expand for multi-donor - donor_device = "cpu" - for dev in sorted_devices: - if dev != compute_device: - donor_device = dev - break # Use first available donor - - block_assignments[block_name] = donor_device - offloaded_bytes += block_memory + if mode == "ratio": + total_requested_percent = sum(raw_parsed.values()) + normalized_ratios = {} + if total_requested_percent > 0: + for dev, percent in raw_parsed.items(): + normalized_ratios[dev] = (percent / total_requested_percent) * 100 else: - # Assign remaining blocks to the primary compute device - block_assignments[block_name] = compute_device + even_share = 100.0 / len(sorted_devices) if sorted_devices else 0 + for dev in sorted_devices: + normalized_ratios[dev] = even_share + + exact_allocations = {} + for dev in sorted_devices: + if dev in normalized_ratios: + exact_allocations[dev] = (normalized_ratios[dev] / 100) * total_memory + else: + exact_allocations[dev] = 0 + + sorted_blocks = sorted(block_list, key=lambda b: b[3], reverse=True) + + device_remaining = exact_allocations.copy() + for block_name, module, block_type, block_memory in sorted_blocks: + best_device = None + max_remaining = -1 + for device in sorted_devices: + if device_remaining[device] >= block_memory and device_remaining[device] > max_remaining: + best_device = device + max_remaining = device_remaining[device] + + if best_device is None: + best_device = compute_device + + block_assignments[block_name] = best_device + device_remaining[best_device] -= block_memory + + unassigned_blocks = [b for b in block_list if b[0] not in block_assignments] + if unassigned_blocks: + unassigned_memory = sum(b[3] for b in unassigned_blocks) + logger.warning(f"[MultiGPU_DisTorch2] {unassigned_memory / (1024**2):.2f} MB of model did not fit into ratio allocations. Assigning to compute device '{compute_device}'.") + for block_name, _, _, _ in unassigned_blocks: + block_assignments[block_name] = compute_device + else: + for block_name, module, block_type, block_memory in reversed(block_list): + for device in devices_to_fill: + if device_quotas.get(device, 0) >= block_memory: + block_assignments[block_name] = device + device_quotas[device] -= block_memory + break - # Explicitly assign tiny blocks to the compute device - if tiny_block_list: - for block_name, module, block_type, block_memory in tiny_block_list: - block_assignments[block_name] = compute_device + unassigned_blocks = [b for b in block_list if b[0] not in block_assignments] + if unassigned_blocks: + if mode == 'gb': + logger.info(f"[MultiGPU_DisTorch2-GB] Assigning {len(unassigned_blocks)} remaining blocks to wildcard device '{wildcard_device}'.") + for block_name, _, _, _ in unassigned_blocks: + block_assignments[block_name] = wildcard_device + else: + unassigned_memory = sum(b[3] for b in unassigned_blocks) + logger.warning(f"[MultiGPU_DisTorch2] {unassigned_memory / (1024**2):.2f} MB of model did not fit into allocations. Assigning to compute device '{compute_device}'.") + for block_name, _, _, _ in unassigned_blocks: + block_assignments[block_name] = compute_device - # Populate device_assignments from the final block_assignments + if mode == 'gb': + total_non_wildcard_quota = sum(v for k, v in parsed_allocations.items() if k != wildcard_device) + distributable_memory = sum(b[3] for b in block_list) + if distributable_memory <= total_non_wildcard_quota: + remaining_quota = total_non_wildcard_quota - distributable_memory + logger.info(f"[MultiGPU_DisTorch2-GB] Underflow: Model fits in non-wildcard devices with {remaining_quota / (1024**2):.2f} MB to spare. Wildcard device '{wildcard_device}' will not be used for distributable blocks.") + + for block_name, _, _, _ in tiny_block_list: + block_assignments[block_name] = compute_device + + # Populate final assignments for logging + device_assignments = defaultdict(list) for block_name, device in block_assignments.items(): - # Find the block in the original list to get all its info for b_name, b_module, b_type, b_mem in all_blocks: if b_name == block_name: device_assignments[device].append((b_name, b_module, b_type, b_mem)) break - # Log final assignments - IDENTICAL FORMAT TO GGML + # Log final assignments logger.info("DisTorch2 Model Final Device/Layer Assignments") logger.info(dash_line) logger.info(fmt_assign.format("Device", "Layers", "Memory (MB)", "% Total")) logger.info(dash_line) - # Calculate and log tiny blocks separately if tiny_block_list: tiny_block_memory = sum(b[3] for b in tiny_block_list) tiny_mem_mb = tiny_block_memory / (1024 * 1024) tiny_mem_percent = (tiny_block_memory / total_memory) * 100 if total_memory > 0 else 0 device_label = f"{compute_device} (<0.01%)" logger.info(fmt_assign.format(device_label, str(len(tiny_block_list)), f"{tiny_mem_mb:.2f}", f"{tiny_mem_percent:.1f}%")) - logger.debug(f"[MultiGPU_DisTorch2] Tiny block memory breakdown: {tiny_block_memory} bytes ({tiny_mem_mb:.2f} MB), which is {tiny_mem_percent:.4f}% of total model memory.") - # Log distributed blocks - total_assigned_memory = 0 device_memories = {} - for device, blocks in device_assignments.items(): - # Exclude tiny blocks from this calculation dist_blocks = [b for b in blocks if b[3] >= MIN_BLOCK_THRESHOLD] - if not dist_blocks: - continue + if not dist_blocks: continue + device_memories[device] = sum(b[3] for b in dist_blocks) - device_memory = sum(b[3] for b in dist_blocks) - device_memories[device] = device_memory - total_assigned_memory += device_memory + final_sorted_devices = sorted(device_memories.keys(), key=lambda d: (d == "cpu", d)) - sorted_assignments = sorted(device_memories.keys(), key=lambda d: (d == "cpu", d)) - - for dev in sorted_assignments: - # Get only the distributed blocks for the count + for dev in final_sorted_devices: dist_blocks = [b for b in device_assignments[dev] if b[3] >= MIN_BLOCK_THRESHOLD] - if not dist_blocks: - continue - + if not dist_blocks: continue mem_mb = device_memories[dev] / (1024 * 1024) mem_percent = (device_memories[dev] / total_memory) * 100 if total_memory > 0 else 0 logger.info(fmt_assign.format(dev, str(len(dist_blocks)), f"{mem_mb:.2f}", f"{mem_percent:.1f}%"))