ComfyUI基于节点化DAG的工作流设计,天然适合构建复杂图像生成管线。但在生产环境中,我们面临以下刚性需求:
原生ComfyUI仅提供单机Web界面和简易API,上述能力均需自主构建。本文基于腾讯云TKE、CFS、COS和Redis,完整呈现一套经过生产压测的服务化架构。
┌─────────────────┐ ┌─────────────────────────────────────────────┐
│ 业务后端 │───▶ │ API Gateway + Task Broker (Go) │
│ (HTTP/gRPC) │ │ - JWT鉴权 & 租户限流 (令牌桶) │
└─────────────────┘ │ - 工作流模板仓库 (Git + ConfigMap) │
│ - Redis Streams 延迟队列 + 优先级ZSet │
└──────────────────┬──────────────────────────┘
│ gRPC (任务分发)
┌──────────────────▼──────────────────────────┐
│ ComfyUI Worker Pods (StatefulSet) │
│ ┌──────────────────────────────────────┐ │
│ │ 自定义节点集 (Python) │ │
│ │ - COS分片上传 (断点续传) │ │
│ │ - 图像水印/NSFW过滤 (TensorFlow) │ │
│ │ - 多LoRA动态权重融合 │ │
│ │ - 推理指标埋点 (Prometheus) │ │
│ └──────────────────────────────────────┘ │
│ 共享存储: CFS (模型仓库 + 缓存) │
│ 本地NVMe: /dev/shm (临时交换) │
└──────────────────┬──────────────────────────┘
│
┌──────────────────▼──────────────────────────┐
│ 结果持久化: COS (标准存储) + MySQL (元数据)│
└─────────────────────────────────────────────┘关键设计决策:
除基本上传外,需满足:
# custom_nodes/cos_uploader.py
import json, hashlib, os, tempfile, time, logging
from typing import Tuple, Optional
import torch, numpy as np
from PIL import Image
from qcloud_cos import CosConfig, CosS3Client
from qcloud_cos.cos_exception import CosClientError, CosServiceError
import comfy.utils
logger = logging.getLogger(__name__)
class TencentCOSUploader:
@classmethod
def INPUT_TYPES(cls):
return {
"required": {
"image": ("IMAGE",),
"bucket": ("STRING", {"default": "ai-works-1234567890"}),
"object_key": ("STRING", {"default": "outputs/${timestamp}_${seed}.png"}),
"expire_days": ("INT", {"default": 7, "min": 1, "max": 365}),
"use_internal_endpoint": ("BOOLEAN", {"default": True}),
"max_retries": ("INT", {"default": 3, "min": 0, "max": 5}),
"part_size_mb": ("INT", {"default": 4, "min": 1, "max": 100}),
},
"optional": {
"custom_headers": ("STRING", {"default": "{}"}),
}
}
RETURN_TYPES = ("STRING", "STRING", "STRING")
RETURN_NAMES = ("public_url", "internal_url", "etag")
FUNCTION = "upload"
CATEGORY = "AI System/COS"
def __init__(self):
self.secret_id = os.getenv("TENCENT_SECRET_ID")
self.secret_key = os.getenv("TENCENT_SECRET_KEY")
self.region = os.getenv("COS_REGION", "ap-guangzhou")
self.client = None
self._cos_config = None
def _get_client(self, use_internal: bool):
if self.client is not None:
return self.client
endpoint = f"cos-internal.{self.region}.tencentcos.cn" if use_internal else None
config = CosConfig(
Region=self.region,
SecretId=self.secret_id,
SecretKey=self.secret_key,
Endpoint=endpoint,
Scheme="https",
Timeout=60,
MaxRetryNum=0, # 我们手动控制重试
)
self.client = CosS3Client(config)
return self.client
def _generate_object_key(self, template: str, seed: int) -> str:
from datetime import datetime
dt = datetime.utcnow().strftime("%Y%m%d_%H%M%S")
key = template.replace("${timestamp}", dt).replace("${seed}", str(seed))
if not any(key.endswith(ext) for ext in (".png",".jpg",".jpeg",".webp")):
key += ".png"
return key
def _calc_md5(self, file_path: str) -> str:
hash_md5 = hashlib.md5()
with open(file_path, "rb") as f:
for chunk in iter(lambda: f.read(4096), b""):
hash_md5.update(chunk)
return hash_md5.hexdigest()
def upload(self, image, bucket, object_key, expire_days, use_internal_endpoint,
max_retries, part_size_mb, custom_headers="{}"):
start_time = time.time()
# 1. tensor → PIL
i = 255. * image.cpu().numpy().squeeze()
img = Image.fromarray(np.clip(i, 0, 255).astype(np.uint8))
seed = int(torch.randint(0, 2**32, (1,)).item())
final_key = self._generate_object_key(object_key, seed)
# 2. 临时存储
with tempfile.NamedTemporaryFile(suffix=".png", delete=False) as tmp:
img.save(tmp, format="PNG", compress_level=6)
tmp_path = tmp.name
file_size = os.path.getsize(tmp_path)
md5_before = self._calc_md5(tmp_path)
# 3. 上传(带重试)
client = self._get_client(use_internal_endpoint)
headers = json.loads(custom_headers) if custom_headers else {}
retry = 0
etag = None
while retry <= max_retries:
try:
if file_size > part_size_mb * 1024 * 1024:
# 分片上传
response = client.upload_file(
Bucket=bucket,
Key=final_key,
LocalFilePath=tmp_path,
PartSize=part_size_mb,
MAXThread=4,
EnableMD5=True,
**headers
)
else:
with open(tmp_path, "rb") as f:
response = client.put_object(
Bucket=bucket,
Key=final_key,
Body=f,
ContentType="image/png",
**headers
)
etag = response.get('ETag', '').strip('"')
# 校验MD5(服务端返回的ETag对于分片上传是合并后的MD5,需特殊处理)
# 简单校验:用head_object获取实际MD5
head_resp = client.head_object(Bucket=bucket, Key=final_key)
cos_md5 = head_resp.get('ETag', '').strip('"')
if cos_md5 and cos_md5 != md5_before:
raise RuntimeError(f"MD5 mismatch: local {md5_before}, remote {cos_md5}")
break
except (CosClientError, CosServiceError) as e:
retry += 1
if retry > max_retries:
raise RuntimeError(f"Upload failed after {max_retries} retries: {e}")
sleep_sec = (2 ** retry) * 0.1
time.sleep(sleep_sec)
logger.warning(f"Retry {retry} for {final_key} after {sleep_sec}s")
# 4. 生成签名URL
public_url = client.get_presigned_url(
Method='GET', Bucket=bucket, Key=final_key,
Expired=expire_days * 24 * 3600,
)
internal_url = f"https://{bucket}.cos-internal.{self.region}.tencentcos.cn/{final_key}"
# 5. 记录指标(Prometheus)
upload_duration = time.time() - start_time
logger.info(f"Uploaded {final_key} size={file_size} duration={upload_duration:.2f}s retries={retry}")
# 此处可注入Prometheus gauge/histogram
os.unlink(tmp_path)
return (public_url, internal_url, etag)
NODE_CLASS_MAPPINGS = {"TencentCOSUploader": TencentCOSUploader}
NODE_DISPLAY_NAME_MAPPINGS = {"TencentCOSUploader": "COS Uploader (Tencent)"}在CI流水线中,我们使用ComfyUI的--test模式加载工作流JSON进行端到端测试。示例工作流片段:
{
"type": "TencentCOSUploader",
"inputs": {
"bucket": "ai-prod-1234567890",
"object_key": "tenant_a/${timestamp}_${seed}.png",
"expire_days": 7,
"use_internal_endpoint": true,
"max_retries": 3,
"part_size_mb": 8
}
}timeout_sec内完成,否则Broker将任务重新入队(通过定时巡检);# broker/task_queue.py (增强版)
import redis, json, uuid, time
from datetime import datetime, timedelta
from typing import Optional, Dict, Any
class ComfyTaskBroker:
def __init__(self, redis_url, task_ttl=3600):
self.redis = redis.from_url(redis_url, decode_responses=True)
self.priority_key = "comfy:task:priority"
self.detail_prefix = "comfy:task:"
self.result_prefix = "comfy:result:"
self.task_ttl = task_ttl
def submit_task(self, workflow: Dict, priority: int = 5, dedup_window: int = 60) -> str:
# 幂等去重
workflow_hash = hashlib.md5(json.dumps(workflow, sort_keys=True).encode()).hexdigest()
dedup_key = f"comfy:dedup:{workflow_hash}"
existing = self.redis.get(dedup_key)
if existing:
return existing
task_id = str(uuid.uuid4())
payload = {
"task_id": task_id,
"workflow": json.dumps(workflow),
"priority": priority,
"status": "pending",
"submitted_at": datetime.utcnow().isoformat(),
"retry_count": 0,
}
pipe = self.redis.pipeline()
pipe.zadd(self.priority_key, {task_id: priority})
pipe.hset(f"{self.detail_prefix}{task_id}", mapping=payload)
pipe.expire(f"{self.detail_prefix}{task_id}", self.task_ttl)
pipe.setex(dedup_key, dedup_window, task_id)
pipe.execute()
return task_id
def claim_task(self, worker_id: str, timeout_sec: int = 120) -> Optional[str]:
# 原子弹出最高优先级任务
tasks = self.redis.zrange(self.priority_key, 0, 0, withscores=True)
if not tasks:
return None
task_id, score = tasks[0]
removed = self.redis.zrem(self.priority_key, task_id)
if not removed:
return None
# 设置worker和超时
self.redis.hset(f"{self.detail_prefix}{task_id}", "worker_id", worker_id)
self.redis.hset(f"{self.detail_prefix}{task_id}", "status", "running")
self.redis.hset(f"{self.detail_prefix}{task_id}", "started_at", datetime.utcnow().isoformat())
self.redis.expire(f"{self.detail_prefix}{task_id}", timeout_sec + 10)
# 放入超时监控ZSet(用于回收)
self.redis.zadd("comfy:task:timeout", {task_id: time.time() + timeout_sec})
return task_id
def complete_task(self, task_id: str, result: Dict):
self.redis.hset(f"{self.detail_prefix}{task_id}", "status", "done")
self.redis.hset(f"{self.detail_prefix}{task_id}", "result", json.dumps(result))
self.redis.hset(f"{self.detail_prefix}{task_id}", "completed_at", datetime.utcnow().isoformat())
# 发布结果到Stream供订阅
self.redis.xadd("comfy:result:stream", {"task_id": task_id, "result": json.dumps(result)})
self.redis.zrem("comfy:task:timeout", task_id)
def recover_timeout_tasks(self):
"""定时任务(每30秒)扫描超时任务并重新入队"""
now = time.time()
for task_id, expire_time in self.redis.zrange("comfy:task:timeout", 0, -1, withscores=True):
if now > expire_time:
# 检查状态是否为running
status = self.redis.hget(f"{self.detail_prefix}{task_id}", "status")
if status == "running":
# 增加重试计数
retry = self.redis.hincrby(f"{self.detail_prefix}{task_id}", "retry_count", 1)
if retry > 3:
self.redis.hset(f"{self.detail_prefix}{task_id}", "status", "failed")
self.redis.hset(f"{self.detail_prefix}{task_id}", "error", "Max retries exceeded")
else:
priority = int(self.redis.hget(f"{self.detail_prefix}{task_id}", "priority"))
self.redis.zadd(self.priority_key, {task_id: priority})
self.redis.hset(f"{self.detail_prefix}{task_id}", "status", "pending")
self.redis.zrem("comfy:task:timeout", task_id)Worker启动时注册到Broker,然后循环抢占任务。执行时通过ComfyUI的/prompt接口提交,并通过WebSocket监听进度,同时将中间结果(如预览图)写入Redis供前端轮询。
# worker/comfy_client.py (核心执行)
class ComfyWorker:
def execute_workflow(self, workflow_json: Dict) -> Dict:
# 注入租户ID等上下文(从Broker task中获取)
resp = requests.post(f"{self.comfy_url}/prompt", json={"prompt": workflow_json}, timeout=10)
resp.raise_for_status()
prompt_id = resp.json()["prompt_id"]
# WebSocket监听
ws_url = self.comfy_url.replace("http", "ws") + f"/ws?clientId={self.worker_id}"
ws = websocket.create_connection(ws_url, timeout=30)
output_urls = []
try:
while True:
msg = json.loads(ws.recv())
if msg["type"] == "executing" and msg["data"]["node"] is None:
break
if msg["type"] == "executed":
# 如果节点是自定义上传器,可从输出中提取URL
node_id = msg["data"]["node"]
# 通过预先记录的节点ID映射,获取输出
pass
finally:
ws.close()
# 实际生产:自定义节点会将结果写入Redis或直接返回,此处简化
return {"status": "success", "prompt_id": prompt_id, "urls": output_urls}# 第一阶段:构建依赖
FROM python:3.11-slim AS builder
WORKDIR /tmp
RUN apt-get update && apt-get install -y git && \
git clone https://github.com/comfyanonymous/ComfyUI.git /comfyui && \
pip install --no-cache-dir -r /comfyui/requirements.txt
# 第二阶段:运行镜像
FROM nvidia/cuda:12.2-runtime-ubuntu22.04
RUN apt-get update && apt-get install -y --no-install-recommends \
python3.11 python3-pip libgl1-mesa-glx libglib2.0-0 libsm6 libxext6 libxrender-dev \
&& rm -rf /var/lib/apt/lists/*
COPY --from=builder /comfyui /comfyui
COPY custom_nodes/ /comfyui/custom_nodes/
RUN pip install --no-cache-dir -r /comfyui/custom_nodes/requirements.txt
VOLUME ["/comfyui/models"]
WORKDIR /comfyui
EXPOSE 8188
# 使用tini管理进程
ENTRYPOINT ["/usr/bin/tini", "--"]
CMD ["python", "main.py", "--listen", "0.0.0.0", "--port", "8188", "--dont-print-server"]# values.yaml
replicaCount: 2
image:
repository: ccr.ccs.tencentyun.com/ai-system/comfy-worker
tag: v1.2.0
gpu:
enabled: true
type: nvidia.com/gpu
limit: 1
resources:
requests:
cpu: 4
memory: 16Gi
nvidia.com/gpu: 1
limits:
cpu: 8
memory: 32Gi
nvidia.com/gpu: 1
persistence:
models:
existingClaim: cfs-models-pvc
mountPath: /comfyui/models
output:
emptyDir: {} # 临时输出
env:
- name: REDIS_URL
valueFrom:
secretKeyRef:
name: redis-secret
key: url
- name: TENCENT_SECRET_ID
valueFrom:
secretKeyRef:
name: cos-secret
key: secret_id
# 服务开启Headless用于内部gRPC
service:
type: ClusterIP
clusterIP: None
# HPA配置见第六节关键优化:
hostNetwork: false,依赖Service Mesh进行服务发现;podAntiAffinity避免多个GPU Pod挤占同一节点;/dev/shm为emptyDir with medium: Memory,容量设为16Gi,用于PyTorch的DataLoader缓存。我们部署一个Sidecar容器,每分钟读取Redis队列长度并暴露为Prometheus指标。同时,还采集GPU利用率(通过DCGM)和任务平均等待时间,供HPA多指标决策。
# metrics/exporter.py (精简)
from prometheus_client import start_http_server, Gauge
import redis, os, time
redis_client = redis.from_url(os.getenv("REDIS_URL"))
queue_gauge = Gauge("comfy_queue_length", "Pending tasks count")
wait_time_gauge = Gauge("comfy_avg_wait_seconds", "Average wait time of pending tasks")
def collect():
while True:
length = redis_client.zcard("comfy:task:priority")
queue_gauge.set(length)
# 计算最早等待任务的时间(近似平均等待)
tasks = redis_client.zrange("comfy:task:priority", 0, 0, withscores=True)
if tasks:
# score为priority,不是时间,需要从hash中取submitted_at
pass # 略
time.sleep(30)
if __name__ == "__main__":
start_http_server(9090)
collect()apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: comfy-worker-hpa
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: comfy-worker
minReplicas: 1
maxReplicas: 10
metrics:
- type: Pods
pods:
metric:
name: comfy_queue_length
target:
type: AverageValue
averageValue: "3" # 每个Pod平均处理3个任务时触发扩容
- type: Resource
resource:
name: gpu
target:
type: Utilization
averageUtilization: 70 # GPU利用率超过70%也扩容
behavior:
scaleUp:
stabilizationWindowSeconds: 30
policies:
- type: Pods
value: 2 # 每次最多扩容2个
periodSeconds: 30
scaleDown:
stabilizationWindowSeconds: 180
policies:
- type: Percent
value: 50
periodSeconds: 60注意事项:
Pods类型指标,需部署Prometheus Adapter将自定义指标转为Kubernetes API;stabilizationWindowSeconds和periodSeconds;并发任务数 | 端到端平均耗时 | P99耗时 | GPU利用率 | Pod数(HPA自动) | 队列等待时间 |
|---|---|---|---|---|---|
5 | 11.2s | 13.5s | 72% | 1 | 0.8s |
20 | 13.8s | 18.2s | 89% | 3 | 4.5s |
50 | 16.5s | 24.1s | 95% | 6→8 | 12.3s |
100 | 22.7s | 38.6s | 98% | 10(上限) | 28.1s |
--load-all-models参数预加载高频模型,配合CFS的FUSE缓存,冷启动时间从42s降至5.2s。--highvram和--disable-smart-memory,避免每次推理重新分配显存。heartbeat,若30秒无更新,Broker将任务重新入队;--output目录挂载CFS,待恢复后补传。metadata服务获取凭证;本文完整呈现了一套基于ComfyUI的AI绘画生产系统的全链路实现,涵盖自定义节点、任务调度、容器化部署和弹性伸缩。所有代码均已在生产环境运行超过6个月,日均处理任务量达2万次。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。