map 按序回结果
(每道题开头都有同一段:上面的池小工具。)
贯穿全条的池小工具(标准库 ThreadPoolExecutor:固定 workers 个线程复用):
pool_sum(fn, args, workers) 池对每个参数跑 fn,交回所有结果之和(和与顺序无关,确定)
pool_map(fn, args, workers) map 按提交顺序回结果(顺序确定)
one_future(fn, x, workers) 提交一个任务,交回它的 Future 的 result()
submit_all(fn, args, workers) 提交全部,按提交顺序 result() 收结果
pool_peak(ntasks, workers) = min(ntasks, workers) 一次最多几个在跑
reuse(ntasks, workers) = (workers, ntasks) 池建几个线程 / 手写要建几个池 workers=2,对 [1,2,3,4] 各算平方,map 回结果:
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
print(",".join(str(r) for r in pool_map(sq, [1, 2, 3, 4], 2)))
全部评论