Add memory parsing and flexible allocation support for safetensor loading

This commit is contained in:
John Pollock
2025-08-25 19:57:55 -05:00
parent 47ed1bed69
commit b6d3403b71
+295 -101
View File
@@ -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}%"))