AI漫剧制作全流程:Dify+ComfyUI自动化工作流与变现指南
2026/9/20 17:25:02
Python 自 1991 年诞生以来,以其简洁优雅的语法、强大的生态系统和“胶水语言”的灵活性,迅速成为 Web 开发、数据科学、人工智能、自动化等领域的核心语言。随着业务规模增长、实时性需求提升,并发编程成为 Python 开发者必须掌握的能力。
你可能已经使用过:
concurrent.futures.ThreadPoolExecutormultiprocessing.Poolasyncio这些工具极大降低了并发编程的门槛,但也让很多开发者忽略了底层原理。
为什么要手写线程池?
更重要的是:
当你能手写线程池时,你对 Python 并发的理解将从“会用”跃升到“精通”。
为了让初学者也能顺利阅读,我们先快速回顾 Python 并发的基础知识。
Python 的线程由操作系统调度,但 CPython 有 GIL(全局解释器锁),导致:
但线程池的设计思想与 GIL 无关,它是通用的并发模型。
线程池的核心思想:
这就是典型的生产者-消费者模型。
条件变量用于:
线程池中:
为了保持文章结构一致,我们插入一个基础示例:
importtimedeftimer(func):defwrapper(*args,**kwargs):start=time.time()result=func(*args,**kwargs)end=time.time()print(f"{func.__name__}花费时间:{end-start:.4f}秒")returnresultreturnwrapper@timerdefcompute_sum(n):returnsum(range(n))print(compute_sum(1000000))我们将从最小可用版本开始,一步步扩展。
我们需要:
importthreadingfromcollectionsimportdequeclassWorkQueue:def__init__(self):self.queue=deque()self.lock=threading.Lock()self.not_empty=threading.Condition(self.lock)defput(self,item):withself.not_empty:self.queue.append(item)self.not_empty.notify()# 通知等待的线程defget(self):withself.not_empty:whilenotself.queue:self.not_empty.wait()# 队列为空,等待returnself.queue.popleft()Condition(self.lock):条件变量绑定锁wait():释放锁并阻塞,直到被 notifynotify():唤醒一个等待线程while not queue:防止虚假唤醒Worker 线程需要:
classWorker(threading.Thread):def__init__(self,work_queue,pool):super().__init__()self.work_queue=work_queue self.pool=pool self.daemon=True# 主线程退出时自动退出defrun(self):whileTrue:task=self.work_queue.get()iftaskisNone:# 收到关闭信号breakfunc,args,kwargs=tasktry:func(*args,**kwargs)exceptExceptionase:print("任务执行异常:",e)线程池需要:
classThreadPool:def__init__(self,num_workers=4):self.work_queue=WorkQueue()self.workers=[]self.num_workers=num_workers self._init_workers()def_init_workers(self):for_inrange(self.num_workers):worker=Worker(self.work_queue,self)worker.start()self.workers.append(worker)defsubmit(self,func,*args,**kwargs):self.work_queue.put((func,args,kwargs))defshutdown(self,wait=True):# 向每个 worker 发送关闭信号for_inself.workers:self.work_queue.put(None)ifwait:forworkerinself.workers:worker.join()下面是完整代码,可直接运行:
importthreadingfromcollectionsimportdequeclassWorkQueue:def__init__(self):self.queue=deque()self.lock=threading.Lock()self.not_empty=threading.Condition(self.lock)defput(self,item):withself.not_empty:self.queue.append(item)self.not_empty.notify()defget(self):withself.not_empty:whilenotself.queue:self.not_empty.wait()returnself.queue.popleft()classWorker(threading.Thread):def__init__(self,work_queue,pool):super().__init__()self.work_queue=work_queue self.pool=pool self.daemon=Truedefrun(self):whileTrue:task=self.work_queue.get()iftaskisNone:breakfunc,args,kwargs=tasktry:func(*args,**kwargs)exceptExceptionase:print("任务执行异常:",e)classThreadPool:def__init__(self,num_workers=4):self.work_queue=WorkQueue()self.workers=[]self.num_workers=num_workers self._init_workers()def_init_workers(self):for_inrange(self.num_workers):worker=Worker(self.work_queue,self)worker.start()self.workers.append(worker)defsubmit(self,func,*args,**kwargs):self.work_queue.put((func,args,kwargs))defshutdown(self,wait=True):for_inself.workers:self.work_queue.put(None)ifwait:forworkerinself.workers:worker.join()importtimedeftask(n):print(f"开始任务{n}")time.sleep(1)print(f"结束任务{n}")pool=ThreadPool(num_workers=3)foriinrange(10):pool.submit(task,i)pool.shutdown()输出示例:
开始任务 0 开始任务 1 开始任务 2 结束任务 0 开始任务 3 ...如果你愿意,我可以继续扩展:
这些都是生产级线程池需要的能力。
importrequestsdeffetch(url):resp=requests.get(url)print(url,len(resp.text))urls=["https://www.python.org","https://www.github.com","https://www.baidu.com",]*3pool=ThreadPool(5)forurlinurls:pool.submit(fetch,url)pool.shutdown()我们从 Python 基础讲起,一步步构建了:
你不仅学会了“如何写”,更理解了“为什么这样写”。
我很想听听你的想法:
告诉我你的方向,我可以继续为你构建更完整的并发体系文章。