手写池跑一批
(每道题开头都有同一段:上面的池小工具。)
本节的手写简单池:
Pool(workers) 起 workers 个 worker 线程,各自 while True 从 queue.Queue 取任务干
p.submit(fn, x) 把 (fn, x) 放进队列
p.shutdown() 放 workers 颗毒丸(None)让 worker 收工,再 join
run_hand(fn, args, workers) → (结果之和, 处理了几个)手写池 workers=3,跑 [1,2,3,4,5] 的平方,看结果之和、处理了几个:
import concurrent.futures
def pool_sum(fn, args, workers):
"""用固定 workers 个线程的池,对每个参数跑 fn,交回所有结果之和(和与顺序无关,确定)。"""
with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as ex:
return sum(ex.map(fn, args))
def pool_map(fn, args, workers):
"""池对每个参数跑 fn,map 按提交顺序回结果(顺序确定)。交回结果列表。"""
with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as ex:
return list(ex.map(fn, args))
def one_future(fn, x, workers):
"""提交一个任务,交回它的 Future 的 result()。"""
with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as ex:
fut = ex.submit(fn, x)
return fut.result()
def submit_all(fn, args, workers):
"""提交全部,按提交顺序 result() 收结果(确定)。交回结果列表。"""
with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as ex:
futs = [ex.submit(fn, a) for a in args]
return [f.result() for f in futs]
def pool_peak(ntasks, workers):
"""池里一次最多几个任务在跑:worker 和任务里少的那个。"""
return min(ntasks, workers)
def reuse(ntasks, workers):
"""跑 ntasks 个任务:线程池只建 workers 个线程复用;手写「一个任务一个线程」要建 ntasks 个。交回 (池建几个, 手写建几个)。"""
return workers, ntasks
import queue
import threading
class Pool:
"""手写的简单线程池:几个 worker 循环从队列取任务、干完接着取;submit 提交,shutdown 放毒丸收工。"""
def __init__(self, workers):
self.workers = workers
self.q = queue.Queue()
self.results = []
self.lock = threading.Lock()
self.threads = [threading.Thread(target=self._work, daemon=True) for _ in range(workers)]
for t in self.threads:
t.start()
def _work(self):
while True:
job = self.q.get()
if job is None: # 毒丸:收工
return
fn, x = job
r = fn(x)
with self.lock:
self.results.append(r)
def submit(self, fn, x):
self.q.put((fn, x))
def shutdown(self):
for _ in range(self.workers):
self.q.put(None)
for t in self.threads:
t.join(timeout=2)
def run_hand(fn, args, workers):
"""用手写池跑一批任务,交回 (结果之和, 一共处理了几个)。"""
p = Pool(workers)
for a in args:
p.submit(fn, a)
p.shutdown()
return sum(p.results), len(p.results)
def due(events, now):
"""events: [(名字, 到点时刻)];交回到 now 为止该跑的任务名(at <= now),按名字排。"""
return sorted(name for name, at in events if at <= now)
def run_order(events):
"""按 (到点时刻, 名字) 排出执行顺序,交回名字列表——定时任务谁先跑是确定的。"""
return [name for name, at in sorted(events, key=lambda e: (e[1], e[0]))]
def every(period, span):
"""一个每 period 跑一次的后台任务,在 [1, span] 时间里跑几次(span // period)。"""
return span // period
def sq(x):
return x * x
total, n = run_hand(sq, [1, 2, 3, 4, 5], 3)
print(str(total) + "/" + str(n))
全部评论