Speedup Gemma4 text generation (CORE-371) (#15054)

This commit is contained in:
Jukka Seppänen 2026-08-14 04:12:55 +03:00 committed by GitHub
parent a779de4d89
commit 7fe8a61385
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
3 changed files with 260 additions and 21 deletions

View File

@ -1,4 +1,5 @@
import torch
import warnings
import weakref
import comfy_aimdo.model_vbar
@ -28,6 +29,18 @@ def cleanup_prefetched_modules(module, comfy_modules):
comfy_aimdo.model_vbar.vbar_unpin(module._v_block)
del module._v_block_faulted
def _drop_graph(module):
graph = getattr(module, "_comfy_graph", None)
if graph is None:
return
# reset() through the bound method surfaces the allocator's benign
# "uncaptured free of a captured allocation" as catchable Python warnings;
# a plain del frees from the C++ dealloc path and spams stderr instead
with warnings.catch_warnings():
warnings.simplefilter("ignore")
graph["graph"].reset()
del module._comfy_graph
def cleanup_prefetch_queues():
global PREFETCH_QUEUES, GRAPH_CAPTURE_STREAMS
@ -41,7 +54,7 @@ def cleanup_prefetch_queues():
cleanup_prefetched_modules(prefetched_module, comfy_modules)
PREFETCH_QUEUES = []
for module in GRAPH_MODULES:
del module._comfy_graph
_drop_graph(module)
GRAPH_MODULES.clear()
GRAPH_WARMED_MODULES.clear()
GRAPH_CAPTURE_STREAMS = {}
@ -117,6 +130,7 @@ def prefetch_queue_pop(queue, device, module, dtype=None, core=None, enable_grap
if signature is not None:
module._v_block_faulted = True
if signature is not None:
_drop_graph(module)
graph = torch.cuda.CUDAGraph()
if generator is not None:
graph.register_generator_state(generator)

View File

@ -1639,7 +1639,7 @@ def mixed_precision_ops(quant_config={}, compute_dtype=torch.bfloat16, full_prec
self.norm_type, self.scale_grad_by_freq, self.sparse)
target_dtype = out_dtype if out_dtype is not None else weight._params.orig_dtype
x = x.to(dtype=target_dtype)
if scale is not None and scale != 1.0:
if scale is not None:
x = x * scale.to(dtype=target_dtype)
return x

View File

@ -9,10 +9,12 @@ import math
from comfy import sd1_clip
import comfy.model_management
import comfy.model_prefetch
import comfy.ops
import comfy.quant_ops
from comfy.ldm.modules.attention import optimized_attention_for_device
from comfy.rmsnorm import rms_norm
from comfy.text_encoders.llama import RMSNorm, MLP, BaseLlama, BaseGenerate, _make_scaled_embedding
from comfy.text_encoders.llama import RMSNorm, MLP, BaseLlama, BaseGenerate, FixedKV, _make_scaled_embedding
# Intentional minor divergences from transformers -reference implementation:
@ -109,7 +111,28 @@ class Gemma4_12B_Config(Gemma4Config):
suppress_tokens = [258883, 258882]
# unfused RoPE as addcmul_ RoPE diverges from reference code
class RingKV(FixedKV):
# sliding-window ring: writes wrap at capacity, validity saturates
def prepare(self, num_tokens):
capacity = self.key.shape[2]
self.position.fill_(self.index % capacity)
self.seqlen.fill_(min(self.index + num_tokens, capacity))
def _fixed_kv_decode_mask(mask, cache, min_val):
capacity = cache.key.shape[2]
valid = min(cache.index + 1, capacity)
output = mask.new_full((*mask.shape[:-1], capacity), min_val)
if isinstance(cache, RingKV):
positions = torch.arange(cache.index + 1 - valid, cache.index + 1, device=mask.device) % capacity
output.index_copy_(-1, positions, mask[..., -valid:])
else:
output[..., :valid] = mask[..., :valid]
return output
# unfused RoPE as addcmul_ RoPE diverges from reference code (vision only; text
# layers use the kitchen split-half kernel, bitwise-equal to this with bf16 freqs)
def _apply_rotary_pos_emb(x, freqs_cis):
cos, sin = freqs_cis[0], freqs_cis[1]
half = x.shape[-1] // 2
@ -140,6 +163,23 @@ class Gemma4Attention(nn.Module):
if config.k_norm == "gemma3":
self.k_norm = RMSNorm(head_dim, eps=config.rms_norm_eps, device=device, dtype=dtype)
def _decode_attention(self, xq, cache, bias):
if bias is None:
# eager decode: slice the cache to the valid length (python-side index,
# no mask needed; a full ring is order-invariant under softmax)
n = min(cache.index + 1, cache.key.shape[2])
gqa_kwargs = {"enable_gqa": True} if self.num_heads != self.num_kv_heads else {}
attention = optimized_attention_for_device(xq.device, mask=False, small_input=True)
return attention(xq, cache.key[:, :, :n], cache.value[:, :, :n], self.num_heads, skip_reshape=True, scale=1.0, **gqa_kwargs)
# graph capture: fixed-length masked attention over the full capacity, explicit
# math (SDPA leaves its fast path on broadcast-bias + GQA and costs ~0.5ms/layer)
batch_size = xq.shape[0]
groups = self.num_heads // self.num_kv_heads
q = xq.reshape(batch_size, self.num_kv_heads, groups, self.head_dim)
scores = q @ cache.key.transpose(-1, -2) + bias
probs = torch.softmax(scores.float(), dim=-1).to(xq.dtype)
return (probs @ cache.value).reshape(batch_size, 1, self.inner_size)
def forward(
self,
hidden_states: torch.Tensor,
@ -156,10 +196,16 @@ class Gemma4Attention(nn.Module):
if self.q_norm is not None:
xq = self.q_norm(xq)
if isinstance(shared_kv, FixedKV):
# decode on a KV-shared layer: attend the source layer's fixed cache
xq = comfy.quant_ops.ck.apply_rope_split_half1(xq, freqs_cis)
output = self._decode_attention(xq, shared_kv, attention_mask)
return self.o_proj(output), None, None
if shared_kv is not None:
xk, xv = shared_kv
# Apply RoPE to Q only (K already has RoPE from source layer)
xq = _apply_rotary_pos_emb(xq, freqs_cis)
xq = comfy.quant_ops.ck.apply_rope_split_half1(xq, freqs_cis)
present_key_value = None
shareable_kv = None
else:
@ -173,11 +219,39 @@ class Gemma4Attention(nn.Module):
xv = rms_norm(xv)
xk = xk.transpose(1, 2)
xv = xv.transpose(1, 2)
xq = _apply_rotary_pos_emb(xq, freqs_cis)
xk = _apply_rotary_pos_emb(xk, freqs_cis)
xq = comfy.quant_ops.ck.apply_rope_split_half1(xq, freqs_cis)
xk = comfy.quant_ops.ck.apply_rope_split_half1(xk, freqs_cis)
present_key_value = None
if past_key_value is not None:
fixed_cache = past_key_value if isinstance(past_key_value, FixedKV) else None
if fixed_cache is not None:
if seq_length == 1:
# CUDA-graphable decode: write at the device-side ring/linear position
fixed_cache.key.index_copy_(2, fixed_cache.position, xk)
fixed_cache.value.index_copy_(2, fixed_cache.position, xv)
output = self._decode_attention(xq, fixed_cache, attention_mask)
return self.o_proj(output), fixed_cache, None
# prefill: attend the local sequence, persist the tail into the cache
capacity = fixed_cache.key.shape[2]
index = fixed_cache.index
if index + seq_length <= capacity:
fixed_cache.key[:, :, index:index + seq_length] = xk
fixed_cache.value[:, :, index:index + seq_length] = xv
if index > 0:
xk = fixed_cache.key[:, :, :index + seq_length]
xv = fixed_cache.value[:, :, :index + seq_length]
elif index == 0:
# prefill longer than the sliding ring: attend the full local K/V
# (per-query windows come from the prefill sliding mask), cache only
# the last `capacity` keys at their wrapped slots (position % capacity)
slots = torch.arange(seq_length - capacity, seq_length, device=xk.device) % capacity
fixed_cache.key.index_copy_(2, slots, xk[:, :, -capacity:])
fixed_cache.value.index_copy_(2, slots, xv[:, :, -capacity:])
else:
raise RuntimeError("gemma4: chunked prefill past the sliding window is not supported")
present_key_value = fixed_cache
elif past_key_value is not None:
cumulative_len = 0
if len(past_key_value) > 0:
past_key, past_value, cumulative_len = past_key_value
@ -245,6 +319,7 @@ class TransformerBlockGemma4(nn.Module):
self.register_buffer("layer_scalar", torch.empty(1, device=device, dtype=dtype))
def forward(self, x, attention_mask=None, freqs_cis=None, past_key_value=None, per_layer_input=None, shared_kv=None):
output = x
sliding_window = None
if self.sliding_attention:
sliding_window = self.sliding_attention
@ -281,7 +356,8 @@ class TransformerBlockGemma4(nn.Module):
x = self.post_per_layer_input_norm(x)
x = residual + x
x = x * comfy.ops.cast_to_input(self.layer_scalar, x)
# in-place into the input buffer so CUDA-graph replays land in the static x
x = torch.mul(x, comfy.ops.cast_to_input(self.layer_scalar, x), out=output)
return x, present_key_value, shareable_kv
@ -290,6 +366,9 @@ class Gemma4Transformer(nn.Module):
def __init__(self, config, device=None, dtype=None, ops=None):
super().__init__()
self.config = config
self.fixed_kv = True
self.prefetch_dynamic_vbars = True
self.graph_dynamic_vbar_blocks = True
self.embed_tokens = _make_scaled_embedding(ops, config.vocab_size, config.hidden_size, config.hidden_size ** 0.5, device, dtype)
@ -298,6 +377,19 @@ class Gemma4Transformer(nn.Module):
for i in range(config.num_hidden_layers)
])
# KV-shared layers never run k_proj/v_proj/k_norm: their never-resolved vbar
# signatures would block layer graph capture, so prefetch only what executes
first_kv_shared = config.num_hidden_layers - config.num_kv_shared_layers if config.num_kv_shared_layers > 0 else config.num_hidden_layers
self._prefetch_units = []
for i, layer in enumerate(self.layers):
if i >= first_kv_shared:
dead = {layer.self_attn.k_proj, layer.self_attn.v_proj, layer.self_attn.k_norm}
self._prefetch_units.append([
m for m in layer.modules() if next(m.children(), None) is None and m not in dead
])
else:
self._prefetch_units.append(layer)
self.norm = RMSNorm(config.hidden_size, eps=config.rms_norm_eps, device=device, dtype=dtype) if config.final_norm else None
# Precompute RoPE inv_freq on CPU to match reference code's exact value
@ -311,6 +403,9 @@ class Gemma4Transformer(nn.Module):
sliding_inv = 1.0 / (config.rope_theta[1] ** (torch.arange(0, config.head_dim, 2).float() / config.head_dim))
self.register_buffer("_sliding_inv_freq", sliding_inv, persistent=False)
if config.suppress_tokens:
self.register_buffer("_suppress_tokens", torch.tensor(config.suppress_tokens, dtype=torch.long), persistent=False)
# Per-layer input mechanism
self.hidden_size_per_layer_input = config.hidden_size_per_layer_input
if self.hidden_size_per_layer_input:
@ -322,19 +417,26 @@ class Gemma4Transformer(nn.Module):
self.hidden_size_per_layer_input, eps=config.rms_norm_eps,
device=device, dtype=dtype)
def get_dynamic_vram__units(self):
return (list(self.layers), []) if self.graph_dynamic_vbar_blocks else ([], [])
def get_past_len(self, past_key_values):
for kv in past_key_values:
if isinstance(kv, FixedKV):
return kv.index
if len(kv) >= 3:
return kv[2]
return 0
def _freqs_from_inv(self, inv_freq, position_ids, device, dtype):
"""Compute cos/sin from stored inv_freq"""
"""Compute per-pair 2x2 rotation matrices [B, 1, S, d/2, 2, 2] from stored inv_freq"""
inv_exp = inv_freq[None, :, None].float().expand(position_ids.shape[0], -1, 1).to(device)
pos_exp = position_ids[:, None, :].float()
freqs = (inv_exp @ pos_exp).transpose(1, 2)
emb = torch.cat((freqs, freqs), dim=-1)
return emb.cos().unsqueeze(1).to(dtype), emb.sin().unsqueeze(1).to(dtype)
cos, sin = freqs.cos(), freqs.sin()
mat = torch.stack((torch.stack((cos, -sin), dim=-1),
torch.stack((sin, cos), dim=-1)), dim=-2)
return mat.unsqueeze(1).to(dtype)
def compute_freqs_cis(self, position_ids, device, dtype=None):
global_freqs = self._freqs_from_inv(self._global_inv_freq, position_ids, device, dtype)
@ -401,6 +503,71 @@ class Gemma4Transformer(nn.Module):
first_kv_shared = self.config.num_hidden_layers - num_kv_shared if num_kv_shared > 0 else self.config.num_hidden_layers
shared_sliding_kv = None # KV from last non-shared sliding layer
shared_global_kv = None # KV from last non-shared global layer
share_source = {}
if num_kv_shared > 0:
for i in range(first_kv_shared):
share_source[bool(self.layers[i].sliding_attention)] = i
prefetch_queue = comfy.model_prefetch.make_prefetch_queue(
list(self._prefetch_units), x.device,
{"prefetch_dynamic_vbars": self.prefetch_dynamic_vbars and past_key_values is not None})
fixed_kv = (past_key_values is not None and len(past_key_values) > 0
and isinstance(past_key_values[0], FixedKV))
decode = fixed_kv and seq_len == 1
# mirror the conditions under which prefetch_queue_pop can actually capture, so
# eager fallbacks keep the sliced decode path instead of the full-capacity one
enable_graph = (decode and mask is None and self.graph_dynamic_vbar_blocks
and prefetch_queue is not None
and hasattr(self.layers[0], "_v_block")
and not comfy.model_management.args.disable_cuda_graphs
and comfy.model_management.is_device_cuda(x.device))
decode_bias = None
decode_masks = None
if decode:
prepared = set()
for kv in past_key_values:
if isinstance(kv, FixedKV) and id(kv.position) not in prepared:
kv.prepare(seq_len)
prepared.add(id(kv.position))
if mask is not None:
decode_masks = {}
for kv in past_key_values:
if isinstance(kv, FixedKV) and id(kv.position) not in decode_masks:
decode_masks[id(kv.position)] = _fixed_kv_decode_mask(mask, kv, min_val)
if enable_graph:
# static buffers + per-capacity attention biases: layer graphs replay against
# stable storage, refreshed eagerly each step
capacities = tuple(sorted({kv.key.shape[2] for kv in past_key_values if isinstance(kv, FixedKV)}))
state_key = (x.shape, x.dtype, x.device, tuple(t.shape for t in freqs_cis), capacities,
None if per_layer_inputs is None else per_layer_inputs.shape)
state = getattr(self, "_comfy_cross_step_state", None)
if state is None or state["key"] != state_key:
state = {"key": state_key,
"x": torch.empty_like(x),
"freqs_cis": [torch.empty_like(t) for t in freqs_cis],
"bias": {c: torch.empty((1, 1, 1, c), dtype=x.dtype, device=x.device) for c in capacities},
"per_layer": None if per_layer_inputs is None else torch.empty_like(per_layer_inputs),
"bias_valid": -1}
self._comfy_cross_step_state = state
comfy.model_management._register_cross_step(self)
state["x"].copy_(x)
for source, target in zip(freqs_cis, state["freqs_cis"]):
target.copy_(source)
x = state["x"]
freqs_cis = state["freqs_cis"]
if per_layer_inputs is not None:
state["per_layer"].copy_(per_layer_inputs)
per_layer_inputs = state["per_layer"]
valid = past_len + 1
for capacity, bias in state["bias"].items():
if state["bias_valid"] != past_len:
bias.fill_(min_val)
bias[..., :min(valid, capacity)] = 0
elif past_len < capacity:
bias[..., past_len:valid] = 0
state["bias_valid"] = valid
decode_bias = state["bias"]
intermediate = None
all_intermediate = None
@ -429,12 +596,36 @@ class Gemma4Transformer(nn.Module):
is_sliding = hasattr(layer, 'sliding_attention') and layer.sliding_attention
if i >= first_kv_shared and num_kv_shared > 0:
shared = shared_sliding_kv if is_sliding else shared_global_kv
if shared is not None:
layer_kwargs['shared_kv'] = shared
if decode:
layer_kwargs['shared_kv'] = past_key_values[share_source[bool(is_sliding)]]
else:
shared = shared_sliding_kv if is_sliding else shared_global_kv
if shared is not None:
layer_kwargs['shared_kv'] = shared
x, current_kv, shareable_kv = layer(x=x, attention_mask=mask, freqs_cis=freqs_cis, past_key_value=past_kv, **layer_kwargs)
if enable_graph:
bias_cache = layer_kwargs.get('shared_kv', past_kv)
layer_mask = decode_bias[bias_cache.key.shape[2]]
elif decode:
bias_cache = layer_kwargs.get('shared_kv', past_kv)
layer_mask = None if decode_masks is None else decode_masks[id(bias_cache.position)]
else:
layer_mask = mask
result = []
def core():
nonlocal x
x, current_kv, shareable_kv = layer(x=x, attention_mask=layer_mask, freqs_cis=freqs_cis, past_key_value=past_kv, **layer_kwargs)
result.append((current_kv, shareable_kv))
comfy.model_prefetch.prefetch_queue_pop(prefetch_queue, x.device, layer, x.dtype, core=core, enable_graph=enable_graph)
if result:
current_kv, shareable_kv = result[0]
else:
# graph replay: the cache already holds this step's write
current_kv, shareable_kv = past_kv, None
next_key_values.append(current_kv if current_kv is not None else ())
# Only track the last sliding/global before the sharing boundary
@ -447,6 +638,14 @@ class Gemma4Transformer(nn.Module):
if i == intermediate_output:
intermediate = x.clone()
if prefetch_queue is not None:
comfy.model_prefetch.prefetch_queue_pop(prefetch_queue, x.device, None)
if fixed_kv:
for kv in past_key_values:
if isinstance(kv, FixedKV):
kv.advance(seq_len)
if self.norm is not None:
x = self.norm(x)
@ -481,14 +680,37 @@ class Gemma4Base(BaseLlama, BaseGenerate, torch.nn.Module):
if cap:
logits = cap * torch.tanh(logits / cap)
if self.model.config.suppress_tokens:
logits[..., self.model.config.suppress_tokens] = torch.finfo(logits.dtype).min
logits.index_fill_(-1, self.model._suppress_tokens, torch.finfo(logits.dtype).min)
return logits
def init_kv_cache(self, batch, max_cache_len, device, execution_dtype):
past_key_values = []
for _ in range(self.model.config.num_hidden_layers):
past_key_values.append(())
return past_key_values
cfg = self.model.config
num_layers = cfg.num_hidden_layers
if not self.model.fixed_kv:
return [() for _ in range(num_layers)]
first_shared = num_layers - cfg.num_kv_shared_layers if cfg.num_kv_shared_layers > 0 else num_layers
# position/seqlen device tensors are shared per cache geometry and filled once per step
trackers = {}
caches = []
for i in range(num_layers):
if i >= first_shared:
caches.append(())
continue
sliding = cfg.sliding_attention[i % len(cfg.sliding_attention)] if cfg.sliding_attention else False
head_dim = cfg.head_dim if sliding else cfg.global_head_dim
k_eq_v = cfg.attention_k_eq_v and not sliding
kv_heads = cfg.num_global_key_value_heads if k_eq_v else cfg.num_key_value_heads
length = min(sliding, max_cache_len) if sliding else max_cache_len
cache_cls = RingKV if sliding else FixedKV
tracker = trackers.get((cache_cls, length))
if tracker is None:
tracker = (torch.empty((1,), device=device, dtype=torch.int64),
torch.empty((batch,), device=device, dtype=torch.int32))
trackers[(cache_cls, length)] = tracker
# zero-init: decode attends full capacity with masked tails, 0*0 stays finite
key = torch.zeros((batch, kv_heads, length, head_dim), device=device, dtype=execution_dtype)
caches.append(cache_cls(key, torch.zeros_like(key), 0, tracker[0], tracker[1]))
return caches
def preprocess_embed(self, embed, device):
if embed["type"] == "image":
@ -1497,6 +1719,9 @@ def gemma4_te(dtype_llama=None, llama_quantization_metadata=None, model_class=No
if dtype_llama is not None:
dtype = dtype_llama
super().__init__(device=device, dtype=dtype, name="gemma4", clip_model=clip_model, model_options=model_options)
def get_dynamic_vram__units(self):
return getattr(self, self.clip).transformer.model.get_dynamic_vram__units()
return Gemma4TEModel_