LiuJuan Z-Image Generator代码实例:异步生成队列+优先级调度实现
LiuJuan Z-Image Generator代码实例异步生成队列优先级调度实现1. 引言如果你用过一些AI图片生成工具可能会遇到这样的烦恼提交一个生成任务后界面就卡住了只能干等着没法同时提交第二个。或者当多个用户同时使用时任务只能一个接一个排队体验很差。今天要介绍的LiuJuan Z-Image Generator就通过一套异步生成队列和优先级调度的代码设计彻底解决了这个问题。它让图片生成变成了“后台任务”你可以连续提交多个提示词然后去喝杯咖啡回来时所有图片都生成好了并且重要的任务还能优先处理。这个工具本身是基于阿里云通义Z-Image扩散模型并加载了LiuJuan自定义的Safetensors权重专门为生成高质量人像和场景图优化。但它的亮点不止于模型更在于这套让生成过程变得丝滑、高效的后台引擎。本文将深入解析这套队列与调度系统的代码实现看它如何让单机工具也能拥有“多任务并行”的体验。2. 项目核心与痛点分析在深入队列系统之前我们先快速了解一下LiuJuan Z-Image Generator的核心以及它为什么要引入异步队列。2.1 工具核心稳定与高效的生成本体这个工具的核心目标是稳定、高质量地生成图片。它做了几件关键事BF16精度优化强制使用torch.bfloat16格式加载模型。这对像RTX 4090这样的显卡更友好能在保证图片质量的同时更好地利用硬件算力。显存碎片治理通过设置max_split_size_mb:128它像一位内存管家把零碎的显存空间整理好大大降低了生成过程中因为显存碎片化而失败的概率。智能权重加载LiuJuan的权重文件可能和原始的Z-Image模型结构不完全一样。工具会自动清洗权重文件里的键名比如去掉多余的transformer.前缀并用一种宽松的模式加载确保自定义权重能成功注入。CPU卸载启动enable_model_cpu_offload()功能把模型暂时不用的部分“请”到CPU内存里待着等需要时再加载回GPU。这能显著降低单次生成任务对GPU显存的占用。2.2 引入异步队列的动机尽管核心生成器已经很高效但一个直接的交互流程是用户输入提示词 - 工具开始生成 - 界面冻结等待 - 生成完成显示图片。这个过程有两个明显痛点阻塞式体验生成期间用户无法进行任何其他操作包括提交新任务。缺乏任务管理如果多人使用或需要批量生成任务只能线性执行无法区分轻重缓急。因此引入一个异步任务队列系统势在必行。它的目标很明确解耦将用户提交请求前端与图片生成计算后端分离开。非阻塞用户提交任务后立即返回生成在后台进行。可管理任务可以排队、可以查看状态、可以设置优先级。可扩展为未来可能的分布式处理打下基础。接下来我们就看看这套系统是如何用代码构建起来的。3. 异步生成队列系统架构整个系统可以看作一个微型的“生产-消费者”模型。Streamlit前端是生产者负责创建任务后台的队列和工作者Worker是消费者负责处理任务。3.1 核心组件与数据流整个系统的运行遵循下图所示的清晰数据流flowchart TD A[用户提交生成请求] -- B[Streamlit前端界面] B -- “创建Task对象br放入PriorityQueue” -- C[任务队列管理器] C -- “Worker线程br监听并获取任务” -- D[异步工作线程] D -- “调用Z-Image生成器br执行生成” -- E[核心生成引擎] E -- “生成完成br更新任务状态” -- F[结果存储与回调] F -- “前端轮询或brWebSocket通知” -- B G[用户界面br实时查看进度与结果] -- B任务Task 封装一个生成请求的所有信息包括唯一ID、提示词、参数、优先级、状态等。它是一个标准的Python数据类dataclass。优先队列PriorityQueue 使用Python标准库queue.PriorityQueue。它不是简单的先进先出而是根据任务的优先级分数priority score来决定谁先出列。分数越低优先级越高。工作者线程Worker Thread 一个或多个在后台运行的线程。它们像不知疲倦的工人持续监听队列。一旦队列中有任务就取出并执行真正的图片生成逻辑。状态管理器 通常用一个字典如task_dict在内存中保存所有任务的状态等待中、生成中、已完成、失败。前端通过任务ID来查询状态。前端交互Streamlit通过轮询定期查询或更高级的st.rerun、会话状态Session State来更新任务状态和显示生成的图片。3.2 任务Task对象设计这是整个系统的基石它定义了“一个生成任务”是什么。import uuid from dataclasses import dataclass, field from enum import Enum from typing import Optional import time class TaskStatus(Enum): 任务状态枚举 PENDING 等待中 RUNNING 生成中 SUCCESS 已完成 FAILED 已失败 dataclass(orderTrue) # orderTrue 用于支持优先队列排序 class GenerationTask: 图片生成任务对象 # 基础信息不参与排序 task_id: str field(compareFalse, default_factorylambda: str(uuid.uuid4())[:8]) prompt: str field(compareFalse) negative_prompt: str field(compareFalse, default) steps: int field(compareFalse, default12) cfg_scale: float field(compareFalse, default2.0) # 队列调度相关参与排序 priority: int field(default5) # 默认优先级1最高10最低 created_at: float field(default_factorytime.time) # 创建时间戳 # 状态与结果 status: TaskStatus field(compareFalse, defaultTaskStatus.PENDING) result_image: Optional[bytes] field(compareFalse, defaultNone) # 生成的图片字节流 error_message: Optional[str] field(compareFalse, defaultNone) # 内部属性用于PriorityQueue排序 # PriorityQueue根据元组第一个元素排序这里用(priority, created_at) # 这样同优先级任务按创建时间先后处理 def __lt__(self, other): return (self.priority, self.created_at) (other.priority, other.created_at)代码解读dataclass(orderTrue)自动生成比较方法使Task对象可排序这是放入PriorityQueue的关键。field(compareFalse)标记这些字段不参与排序比较只用于存储数据。priority和created_at这两个字段共同决定了任务在队列中的顺序。priority值越小越优先如果priority相同则created_at更早值更小的任务优先。__lt__方法自定义了“小于”比较逻辑确保优先队列能正确工作。4. 优先级调度与队列管理器实现有了Task对象我们需要一个“调度中心”来管理它们的生命周期。4.1 队列管理器QueueManager这个类负责接收任务、管理队列、启动工作线程。import threading import queue from typing import Dict from .task import GenerationTask, TaskStatus class GenerationQueueManager: 生成队列管理器单例模式 _instance None _lock threading.Lock() def __new__(cls): with cls._lock: if cls._instance is None: cls._instance super().__new__(cls) cls._instance._initialized False return cls._instance def __init__(self): if self._initialized: return self._initialized True # 核心队列与存储 self.task_queue queue.PriorityQueue(maxsize100) # 设置最大容量防止内存溢出 self.task_dict: Dict[str, GenerationTask] {} # 任务ID到任务的映射 self.dict_lock threading.Lock() # 保护task_dict的锁 # 工作线程 self.worker_thread None self.is_running False # 启动工作线程 self.start_worker() def start_worker(self): 启动后台工作线程 if self.worker_thread is None or not self.worker_thread.is_alive(): self.is_running True self.worker_thread threading.Thread( targetself._worker_loop, daemonTrue, # 设置为守护线程主程序退出时自动结束 nameZ-Image-Generation-Worker ) self.worker_thread.start() print(队列工作线程已启动) def submit_task( self, prompt: str, negative_prompt: str , steps: int 12, cfg_scale: float 2.0, priority: int 5 ) - str: 提交一个新的生成任务 Args: prompt: 正面提示词 negative_prompt: 负面提示词 steps: 迭代步数 cfg_scale: 引导系数 priority: 优先级 (1-10, 1最高) Returns: 任务ID # 创建任务对象 task GenerationTask( promptprompt, negative_promptnegative_prompt, stepssteps, cfg_scalecfg_scale, prioritymax(1, min(10, priority)) # 限制优先级在1-10范围内 ) # 存入任务字典 with self.dict_lock: self.task_dict[task.task_id] task # 放入优先队列 try: self.task_queue.put(task, blockTrue, timeout5) print(f任务已提交: ID{task.task_id}, 优先级{task.priority}, 提示词{prompt[:30]}...) return task.task_id except queue.Full: # 队列已满更新任务状态为失败 with self.dict_lock: task.status TaskStatus.FAILED task.error_message 任务队列已满请稍后重试 return task.task_id def get_task_status(self, task_id: str) - Optional[GenerationTask]: 根据任务ID获取任务状态 with self.dict_lock: return self.task_dict.get(task_id) def _worker_loop(self): 工作线程的主循环 # 这里是简化版实际需要导入你的Z-Image生成器 # from .z_image_generator import ZImageGenerator # generator ZImageGenerator() # 应全局初始化一次 print(工作线程开始运行...) while self.is_running: try: # 从队列获取任务阻塞式队列为空时等待 task: GenerationTask self.task_queue.get(blockTrue, timeout1) # 更新任务状态 with self.dict_lock: task.status TaskStatus.RUNNING print(f开始处理任务: ID{task.task_id}) # 实际调用生成逻辑这里用伪代码 try: # 真实调用示例需替换: # image_bytes generator.generate( # prompttask.prompt, # negative_prompttask.negative_prompt, # num_inference_stepstask.steps, # guidance_scaletask.cfg_scale # ) # 模拟生成过程 import time time.sleep(3) # 模拟生成耗时 image_bytes bfake_image_data # 模拟生成的图片数据 # 更新任务结果 with self.dict_lock: task.status TaskStatus.SUCCESS task.result_image image_bytes print(f任务完成: ID{task.task_id}) except Exception as e: # 生成失败 with self.dict_lock: task.status TaskStatus.FAILED task.error_message str(e) print(f任务失败: ID{task.task_id}, 错误{e}) finally: # 标记任务已完成用于queue.task_done()如果使用JoinableQueue self.task_queue.task_done() except queue.Empty: # 队列为空短暂空闲继续循环 continue except Exception as e: print(f工作线程发生错误: {e}) import time time.sleep(5) # 发生错误时等待一段时间再继续 def shutdown(self): 关闭工作线程 self.is_running False if self.worker_thread: self.worker_thread.join(timeout5) print(队列工作线程已关闭)关键设计解析单例模式确保整个应用中只有一个队列管理器实例避免状态混乱。线程安全使用threading.Lock保护共享的task_dict防止多个线程同时修改导致数据错误。优先级队列queue.PriorityQueue根据Task的__lt__方法自动排序priority值小的任务先出队。守护线程工作线程设置为daemonTrue这样即使它还在运行主程序退出时Python解释器也会直接退出避免僵尸线程。阻塞获取task_queue.get(blockTrue, timeout1)让工作线程在队列空时等待1秒既避免忙等待busy-waiting消耗CPU又能响应关闭信号。4.2 在Streamlit前端集成后端队列系统准备好后需要在Streamlit界面中调用它。import streamlit as st import time from queue_manager import GenerationQueueManager # 初始化队列管理器单例全局唯一 st.cache_resource def get_queue_manager(): return GenerationQueueManager() queue_mgr get_queue_manager() # Streamlit页面标题 st.title(LiuJuan Z-Image Generator - 异步队列版) # 侧边栏任务提交表单 with st.sidebar: st.header(提交新任务) prompt st.text_area(提示词, height100, placeholder描述你想要生成的图片内容...) negative_prompt st.text_area(负面提示词, height60, placeholder描述你不希望出现在图片中的内容...) col1, col2 st.columns(2) with col1: steps st.slider(迭代步数, min_value5, max_value30, value12) with col2: cfg_scale st.slider(CFG Scale, min_value1.0, max_value10.0, value2.0, step0.5) priority st.select_slider( 任务优先级, options[1, 2, 3, 4, 5, 6, 7, 8, 9, 10], value5, help1为最高优先级10为最低优先级 ) if st.button(提交生成任务, typeprimary, use_container_widthTrue): if prompt.strip(): # 提交任务到队列 task_id queue_mgr.submit_task( promptprompt, negative_promptnegative_prompt, stepssteps, cfg_scalecfg_scale, prioritypriority ) st.success(f任务已提交任务ID: {task_id}) st.info(任务已进入后台队列处理请在下方的任务状态中查看进度。) else: st.warning(请输入提示词) # 主区域任务状态监控与结果展示 st.header(任务状态监控) # 使用Session State存储当前查看的任务ID if selected_task_id not in st.session_state: st.session_state.selected_task_id None # 获取所有任务并显示 tasks queue_mgr.task_dict.values() if tasks: # 创建任务状态表格 st.subheader(任务列表) # 按状态分组显示 for status_group in [生成中, 等待中, 已完成, 已失败]: group_tasks [t for t in tasks if t.status.value status_group] if group_tasks: st.caption(f**{status_group}** ({len(group_tasks)}个)) for task in sorted(group_tasks, keylambda x: x.created_at): col1, col2, col3, col4 st.columns([3, 2, 2, 2]) with col1: # 显示简略提示词 short_prompt task.prompt[:40] ... if len(task.prompt) 40 else task.prompt st.text(fID:{task.task_id} - {short_prompt}) with col2: st.text(f优先级: {task.priority}) with col3: status_color { 等待中: blue, 生成中: orange, 已完成: green, 已失败: red }.get(task.status.value, gray) st.markdown(f:{status_color}[{task.status.value}]) with col4: if st.button(查看, keyfview_{task.task_id}): st.session_state.selected_task_id task.task_id # 显示选中的任务详情 if st.session_state.selected_task_id: st.divider() selected_task queue_mgr.get_task_status(st.session_state.selected_task_id) if selected_task: st.subheader(f任务详情 - ID: {selected_task.task_id}) col1, col2 st.columns(2) with col1: st.metric(状态, selected_task.status.value) st.metric(优先级, selected_task.priority) with col2: st.metric(创建时间, time.strftime(%H:%M:%S, time.localtime(selected_task.created_at))) st.metric(迭代步数, selected_task.steps) st.text_area(完整提示词, selected_task.prompt, height100, disabledTrue) st.text_area(负面提示词, selected_task.negative_prompt, height60, disabledTrue) # 显示生成结果或错误信息 if selected_task.status TaskStatus.SUCCESS and selected_task.result_image: st.image(selected_task.result_image, caption生成结果, use_column_widthTrue) # 可以添加下载按钮 st.download_button( label下载图片, dataselected_task.result_image, file_namefgenerated_{selected_task.task_id}.png, mimeimage/png ) elif selected_task.status TaskStatus.FAILED: st.error(f任务失败: {selected_task.error_message}) if st.button(返回任务列表): st.session_state.selected_task_id None st.rerun() else: st.info(暂无任务请在左侧提交新的生成任务。) # 自动刷新每3秒刷新一次任务状态 time.sleep(3) st.rerun()前端设计亮点状态持久化使用st.cache_resource缓存队列管理器确保页面重载时任务不丢失。实时监控通过time.sleep(3)和st.rerun()实现简单的自动刷新让用户能看到任务状态变化。对于更复杂的应用可以考虑WebSocket。清晰的任务分组按状态生成中、等待中、已完成、已失败分组显示任务一目了然。优先级可视化在任务列表中直接显示优先级数字让用户理解调度顺序。5. 高级特性与优化建议上面的代码实现了一个可用的基础系统。在实际项目中还可以考虑以下增强5.1 多工作线程与负载均衡单个工作线程可能成为瓶颈。可以轻松扩展为线程池import concurrent.futures class AdvancedQueueManager(GenerationQueueManager): def __init__(self, num_workers2): super().__init__() self.num_workers num_workers self.worker_threads [] def start_worker(self): self.is_running True for i in range(self.num_workers): thread threading.Thread( targetself._worker_loop, daemonTrue, namefZ-Image-Worker-{i} ) thread.start() self.worker_threads.append(thread) print(f已启动 {self.num_workers} 个工作线程)5.2 任务取消与超时控制允许用户取消排队中的任务并为执行中的任务设置超时def cancel_task(self, task_id: str) - bool: 取消一个尚未开始的任务 with self.dict_lock: task self.task_dict.get(task_id) if not task: return False if task.status TaskStatus.PENDING: # 标记为取消实际从队列移除较复杂需额外设计 task.status TaskStatus.FAILED task.error_message 用户取消 return True return False # 在工作线程中增加超时控制 import signal class TimeoutException(Exception): pass def timeout_handler(signum, frame): raise TimeoutException(生成超时) def _worker_loop(self): while self.is_running: try: task self.task_queue.get(blockTrue, timeout1) # 设置超时例如60秒 signal.signal(signal.SIGALRM, timeout_handler) signal.alarm(60) try: # 执行生成... pass except TimeoutException: task.status TaskStatus.FAILED task.error_message 生成超时 finally: signal.alarm(0) # 取消闹钟5.3 任务进度反馈让用户能看到生成进度而不是简单的“生成中”# 在Task类中添加进度字段 dataclass class GenerationTask: # ... 其他字段 ... progress: float field(compareFalse, default0.0) # 0.0 ~ 1.0 # 在生成器中更新进度 def generate_with_progress(self, prompt, callbackNone): for i in range(total_steps): # 执行一步生成... if callback: callback(i / total_steps) # 回调更新进度 # 工作线程中 def _worker_loop(self): def update_progress(progress): with self.dict_lock: task.progress progress # 调用生成器时传入回调 image_bytes generator.generate_with_progress( prompttask.prompt, progress_callbackupdate_progress )5.4 结果持久化与历史记录将完成的任务保存到数据库或文件系统支持历史查询import json import sqlite3 from datetime import datetime class TaskHistoryManager: def __init__(self, db_pathtasks.db): self.conn sqlite3.connect(db_path, check_same_threadFalse) self._create_table() def _create_table(self): self.conn.execute( CREATE TABLE IF NOT EXISTS generation_tasks ( task_id TEXT PRIMARY KEY, prompt TEXT, negative_prompt TEXT, steps INTEGER, cfg_scale REAL, priority INTEGER, status TEXT, created_at REAL, finished_at REAL, image_path TEXT, error_message TEXT ) ) def save_task(self, task: GenerationTask): # 保存到数据库... pass def load_history(self, limit50): # 从数据库加载历史任务... pass6. 总结通过为LiuJuan Z-Image Generator实现异步生成队列和优先级调度我们成功地将一个阻塞式的单任务工具转变为一个流畅的多任务处理系统。这套系统的价值在于提升用户体验用户提交任务后无需等待可以连续提交多个请求并实时查看所有任务的状态和进度。优化资源利用通过优先级调度重要的任务可以优先处理。后台工作线程可以持续运行避免频繁的模型加载和卸载开销。增强系统健壮性任务队列作为缓冲区可以应对短时间内的请求高峰。失败的任务不会影响整个系统错误信息也能被妥善记录和展示。架构可扩展性这套“生产-消费者”模型为未来扩展打下了基础。例如可以很容易地将工作线程替换为独立的进程甚至是将任务分发到多台机器的分布式工作节点。实现要点回顾**使用PriorityQueue**实现基于优先级的任务调度。通过线程和锁安全地管理共享状态。设计清晰的Task数据类封装任务所有信息。在Streamlit中结合Session State和自动刷新实现准实时的状态更新。虽然本文展示的是集成在单机工具中的方案但它的设计思想同样适用于更复杂的云服务或分布式系统。下次当你需要为任何耗时的AI任务不仅是图片生成也包括文本生成、视频处理等构建交互界面时不妨考虑引入这样一个异步队列系统它会让你的应用显得更加专业和高效。获取更多AI镜像想探索更多AI镜像和应用场景访问 CSDN星图镜像广场提供丰富的预置镜像覆盖大模型推理、图像生成、视频生成、模型微调等多个领域支持一键部署。