本文聚焦于ComfyUI在生产环境中的工程化落地,围绕自定义节点规范、工作流DAG调度、容器镜像优化、TKE Serverless弹性扩缩容以及全链路可观测性五大核心模块,提供一套完整的代码级实现方案。所有组件均基于腾讯云原生服务(SCF、TKE Serverless、CFS、CLS)构建,兼顾高性能与成本效益。
ComfyUI作为Stable Diffusion工作流引擎,已在众多企业中承担核心推理任务。然而,从“原型验证”走向“线上服务”,通常需要解决以下问题:
针对上述痛点,本指南提出一套完整的工程化体系:通过强制节点接口规范保障代码质量,借助腾讯云Serverless容器实现GPU按需使用,并利用Prometheus自定义指标+HPA完成自动扩缩,最终交付一个可运维、可计量、可快速迭代的AI绘画服务。
┌─────────────────────────────────────────────────────────────┐
│ API Gateway │
│ (路由: /infer, /status, /health) │
└───────────────────────────┬─────────────────────────────────┘
│
┌───────────────────────────▼─────────────────────────────────┐
│ 腾讯云函数(SCF)– 调度层 │
│ - 请求入队(Redis Streams) │
│ - 实例健康探针 │
│ - 自定义节点版本校验 │
└───────────────────────────┬─────────────────────────────────┘
│
┌───────────────────────────▼─────────────────────────────────┐
│ TKE Serverless 超级节点(GPU Pod) │
│ ┌────────────────────────────────────────────────────┐ │
│ │ ComfyUI Runner(定制容器) │ │
│ │ - 集成prometheus_client指标采集 │ │
│ │ - 异步DAG执行引擎(基于asyncio) │ │
│ │ - 自定义节点沙箱(PyInstaller隔离) │ │
│ └────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────┘选型说明:
为保证生产环境稳定性,所有自定义节点必须实现以下抽象基类,并编写配套单元测试。
# node_interface.py
from abc import ABC, abstractmethod
from typing import Dict, Any, Optional
import hashlib
class ComfyUINodeBase(ABC):
@property
@abstractmethod
def node_id(self) -> str:
"""唯一标识,如 'CustomUpscaler.v2'"""
pass
@abstractmethod
async def execute(self, inputs: Dict[str, Any], cache: Optional[Dict] = None) -> Dict[str, Any]:
"""异步执行,必须返回 {'output': ..., 'metrics': {...}}"""
pass
@abstractmethod
def validate_input(self, inputs: Dict[str, Any]) -> bool:
"""强校验,失败时抛出 ValidationError"""
pass
def get_checksum(self) -> str:
"""用于版本追踪"""
return hashlib.sha256(self.node_id.encode() + b'::' + self.__version__.encode()).hexdigest()该节点将生成图片直接上传至COS,避免大图通过Base64回传调度层,降低网络开销。
# nodes/cos_uploader.py
import asyncio
from uuid import uuid4
from qcloud_cos import CosConfig, CosS3Client
from comfyui_sdk import ImageTensor # 假设的SDK类型
class COSUploader(ComfyUINodeBase):
def __init__(self, secret_id: str, secret_key: str, region: str, bucket: str):
self.config = CosConfig(Region=region, SecretId=secret_id, SecretKey=secret_key)
self.client = CosS3Client(self.config)
self.bucket = bucket
@property
def node_id(self) -> str:
return "COSUploader.v1"
async def execute(self, inputs, cache=None):
tensor: ImageTensor = inputs.get('image')
if not tensor:
raise ValueError("Missing 'image' input")
png_bytes = await asyncio.to_thread(tensor.save_to_bytes, format='PNG')
key = f"outputs/{inputs.get('prefix', '')}/{uuid4()}.png"
response = await asyncio.to_thread(
self.client.put_object,
Bucket=self.bucket,
Body=png_bytes,
Key=key,
StorageClass='STANDARD',
EnableMD5=False
)
return {
"url": f"https://{self.bucket}.cos.{self.config.Region}.myqcloud.com/{key}",
"etag": response.get('ETag'),
"bucket": self.bucket
}# test_nodes/test_cos_uploader.py
import pytest
from unittest.mock import patch
from nodes.cos_uploader import COSUploader
@pytest.mark.asyncio
async def test_cos_uploader_execute():
uploader = COSUploader("test_id", "test_key", "ap-guangzhou", "mock-bucket")
with patch.object(uploader.client, 'put_object', return_value={'ETag': '"abc123"'}):
result = await uploader.execute({
"image": MockImageTensor(),
"prefix": "test"
})
assert "url" in result
assert "mock-bucket" in result["url"]
assert result["etag"] == '"abc123"'ComfyUI原生JSON仅描述节点拓扑,缺少并发控制。我们基于Python的graphlib和asyncio实现带并发上限的DAG调度器。
# scheduler/dag_scheduler.py
import asyncio
from graphlib import TopologicalSorter
from typing import Dict, List, Set
import time
class ComfyDAGScheduler:
def __init__(self, workflow_json: dict):
self.nodes = workflow_json['nodes']
self.links = workflow_json.get('links', [])
self._build_graph()
def _build_graph(self):
self.graph = {node['id']: set() for node in self.nodes}
for link in self.links:
src_id, src_slot, dst_id, dst_slot = link
self.graph[dst_id].add(src_id)
async def run(self, node_executor_map: Dict[str, callable], max_concurrency: int = 4):
ts = TopologicalSorter(self.graph)
ts.prepare()
semaphore = asyncio.Semaphore(max_concurrency)
results = {}
tasks = []
async def safe_execute(node_id):
async with semaphore:
parent_ids = self.graph[node_id]
inputs = {p: results.get(p) for p in parent_ids}
start = time.perf_counter()
output = await node_executor_map[node_id](inputs)
results[node_id] = output
self._record_metrics(node_id, time.perf_counter() - start)
while ts.is_active():
ready = ts.get_ready()
for node_id in ready:
task = asyncio.create_task(safe_execute(node_id))
tasks.append(task)
done, _ = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)
for t in done:
tasks.remove(t)
ts.done(*ready)
return resultsSCF集成方式:当HTTP请求到达时,SCF将工作流JSON写入Redis Stream并立即返回task_id,TKE Pod异步消费,避免SCF超时(最大900s)。
TKE Serverless Pod冷启动的主要瓶颈在于模型权重(~5GB)和PyTorch依赖(~2GB)。我们采用分层镜像 + CFS共享挂载策略。
FROM nvidia/cuda:12.2.0-runtime-ubuntu22.04 AS base
RUN apt update && apt install -y python3.10 python3-pip && rm -rf /var/lib/apt/lists/*
FROM base AS deps
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt --extra-index-url https://download.pytorch.org/whl/cu121
FROM deps AS code
COPY ./comfyui/ /app/comfyui/
COPY ./custom_nodes/ /app/custom_nodes/
COPY ./scheduler/ /app/scheduler/
WORKDIR /app
ENTRYPOINT ["python", "runner.py"]# runner.py
import os
import asyncio
import time
from pathlib import Path
from comfyui import ComfyUIApp
MODEL_ROOT = os.getenv("MODEL_PATH", "/mnt/models")
CHECKPOINT_DIR = Path(MODEL_ROOT) / "checkpoints"
async def warmup():
from comfyui.ckpt_loader import lazy_load_checkpoint
lazy_load_checkpoint(CHECKPOINT_DIR / "sd_xl_base.safetensors")
print(f"Warmup done in {time.perf_counter() - start:.2f}s")
if __name__ == "__main__":
start = time.perf_counter()
asyncio.run(warmup())
from aiohttp import web
app = web.Application()
# 添加 /infer, /health 路由
web.run_app(app, host="0.0.0.0", port=7860)冷启动实测数据(腾讯云TKE Serverless,16GB显存):
通过prometheus_client暴露Pod的GPU利用率和队列深度,结合腾讯云TKE的自定义指标HPA实现精准扩缩。
# metrics/exporter.py
from prometheus_client import Gauge, start_http_server
import pynvml
import redis
import time
class GPUMetrics:
def __init__(self, redis_client):
pynvml.nvmlInit()
self.handle = pynvml.nvmlDeviceGetHandleByIndex(0)
self.gpu_util = Gauge('comfyui_gpu_util_percent', 'GPU utilization')
self.mem_used = Gauge('comfyui_gpu_mem_used_mb', 'GPU memory used MB')
self.queue_depth = Gauge('comfyui_queue_depth', 'Pending tasks')
self.redis = redis_client
def update(self):
util = pynvml.nvmlDeviceGetUtilizationRates(self.handle)
self.gpu_util.set(util.gpu)
mem = pynvml.nvmlDeviceGetMemoryInfo(self.handle)
self.mem_used.set(mem.used // 1024 // 1024)
self.queue_depth.set(self.redis.llen('comfy_task_queue'))
if __name__ == "__main__":
start_http_server(8000)
redis_client = redis.Redis(host=os.getenv('REDIS_ADDR'), decode_responses=True)
metrics = GPUMetrics(redis_client)
while True:
metrics.update()
time.sleep(5)apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: comfyui-runner-hpa
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: comfyui-runner
minReplicas: 0 # 闲时缩到0,节省成本
maxReplicas: 10
metrics:
- type: Pods
pods:
metric:
name: comfyui_queue_depth
target:
type: AverageValue
averageValue: "2"
- type: Pods
pods:
metric:
name: comfyui_gpu_util_percent
target:
type: AverageValue
averageValue: "60"为每条推理请求生成分布式Trace,并上报至腾讯云日志服务(CLS)。
# tracer.py
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.instrumentation.aiohttp_client import AioHttpClientInstrumentor
provider = TracerProvider()
processor = BatchSpanProcessor(OTLPSpanExporter(endpoint="http://otel-collector:4317"))
provider.add_span_processor(processor)
trace.set_tracer_provider(provider)
async def handle_infer(request):
tracer = trace.get_tracer(__name__)
with tracer.start_as_current_span("comfyui_full_pipeline") as span:
span.set_attribute("workflow_id", request.workflow_id)
span.set_attribute("user_id", request.user_id)
result = await dag_scheduler.run(...)
span.set_attribute("total_latency_ms", int((time.perf_counter()-start)*1000))
# 日志自动同步至CLS(通过Logtail)以下脚本演示通过SDK创建SCF函数与TKE集群(实际生产建议配合Terraform或Helm)。
# deploy.py
from tencentcloud.common import credential
from tencentcloud.scf.v20180416 import scf_client, models as scf_models
from tencentcloud.tke.v20180525 import tke_client, models as tke_models
cred = credential.Credential("YOUR_SECRET_ID", "YOUR_SECRET_KEY")
# 1. 创建SCF调度函数
scf_cli = scf_client.ScfClient(cred, "ap-guangzhou")
req = scf_models.CreateFunctionRequest()
req.FunctionName = "comfyui-scheduler"
req.Code = scf_models.Code(ZipFile=open("./scheduler.zip", "rb").read())
req.Handler = "index.main_handler"
req.Runtime = "Python3.10"
req.Timeout = 30
req.Environment = scf_models.Environment(
Variables=[{"Key": "REDIS_ADDR", "Value": "10.0.0.5:6379"}]
)
resp = scf_cli.CreateFunction(req)
print("SCF created:", resp.FunctionId)
# 2. 创建TKE Serverless集群(示例,实际参数需按需调整)
tke_cli = tke_client.TkeClient(cred, "ap-guangzhou")
cluster_req = tke_models.CreateClusterRequest()
cluster_req.ClusterType = "MANAGED_CLUSTER"
cluster_req.ClusterCIDRSettings = tke_models.ClusterCIDRSettings(ClusterCIDR="172.16.0.0/16")
# 后续使用kubectl或SDK部署Deployment+HPA使用腾讯云CLS压测工具(100并发,持续10分钟,SDXL 1024x1024,20步):
指标 | 数值 |
|---|---|
平均推理耗时 | 12.3s |
P95耗时 | 18.7s(含冷重启) |
最大Pod数 | 6(自动扩缩) |
空闲时Pod数 | 0 |
单次推理成本(含GPU+网络) | ¥0.028 |
相比自建固定GPU集群(如A10包月),本方案在低负载时成本降低约73%。
本指南提供的不是一次性配置,而是一套可持续迭代的工程框架:
当新模型发布时,仅需更新CFS中的权重文件并重启Pod;当自定义节点出现Bug,通过回滚镜像版本即可快速恢复。整个体系可平滑集成至已有CI/CD流水线,为企业级AI绘画服务提供坚实底座。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。