线程池批量巡检
concurrent.futures.ThreadPoolExecutor 是运维批量任务的首选:自动管理线程生命周期、限制最大并发、用 submit / map 提交任务并收集结果。典型场景:批量 HTTP 探活、SSH 端口检测、并行调用只读 API。
基本用法
#!/usr/bin/env python3
"""线程池并行 HTTP 探活,输出可达性与响应时间。"""
import time
from concurrent.futures import ThreadPoolExecutor, as_completed, wait
import httpx
def probe_url(url: str, timeout: float = 3.0) -> dict:
"""单 URL 探活,返回结构化结果。"""
start = time.perf_counter()
try:
resp = httpx.get(url, timeout=timeout, follow_redirects=True)
elapsed = time.perf_counter() - start
return {
"url": url,
"ok": resp.status_code < 400,
"status_code": resp.status_code,
"elapsed_ms": round(elapsed * 1000, 1),
"error": None,
}
except httpx.HTTPError as e:
elapsed = time.perf_counter() - start
return {
"url": url,
"ok": False,
"status_code": None,
"elapsed_ms": round(elapsed * 1000, 1),
"error": str(e),
}
def main() -> None:
urls = [
"https://gitlab.example.com",
"https://harbor.example.com",
"https://nexus.example.com",
"https://grafana.example.com",
]
# max_workers 控制并发,避免打满网络或对端限流
with ThreadPoolExecutor(max_workers=8) as pool:
futures = {pool.submit(probe_url, u): u for u in urls}
for fut in as_completed(futures):
result = fut.result() # 异常会在 result() 时抛出
flag = "OK" if result["ok"] else "FAIL"
print(f"[{flag}] {result['url']} {result['elapsed_ms']}ms err={result['error']}")
if __name__ == "__main__":
main()
依赖:uv add httpx
pool.map(fn, iterable) 保序返回;submit + as_completed 适合先完成先处理(实时进度)。
异常隔离与汇总
from concurrent.futures import ThreadPoolExecutor, as_completed
def risky_task(item: str) -> str:
if item == "bad":
raise ValueError("模拟失败")
return f"ok:{item}"
def batch_with_error_handling(items: list[str]) -> None:
ok_list, err_list = [], []
with ThreadPoolExecutor(max_workers=4) as pool:
future_map = {pool.submit(risky_task, x): x for x in items}
for fut in as_completed(future_map):
item = future_map[fut]
try:
ok_list.append(fut.result())
except Exception as e:
err_list.append({"item": item, "error": str(e)})
print(f"成功 {len(ok_list)},失败 {len(err_list)}")
for e in err_list:
print(f" - {e['item']}: {e['error']}")
单个任务失败不应拖垮整批巡检;务必在 as_completed 循环内 try/except。
运维实践建议
# 推荐:并发数与目标规模匹配
MAX_WORKERS = min(32, len(hosts), 10) # 示例:不超过 10 路 SSH
# 推荐:总体等待时间 + 单任务网络超时同时存在
# probe_url 内仍必须传 httpx timeout;线程不能被 Python 安全地强制杀死。
pool = ThreadPoolExecutor(max_workers=MAX_WORKERS)
futures = {pool.submit(probe_url, u, timeout=5.0): u for u in urls}
done, not_done = wait(futures, timeout=120)
try:
for fut in done:
url = futures[fut]
try:
print(url, fut.result())
except Exception as exc:
print(url, {"ok": False, "error": str(exc)})
for fut in not_done:
# 只能取消尚未开始的任务;正在运行的 I/O 依赖其自身 timeout 返回。
fut.cancel()
print(futures[fut], {"ok": False, "error": "batch deadline exceeded"})
finally:
# 不等待已经运行的线程;cancel_futures 取消仍在队列中的工作项。
pool.shutdown(wait=False, cancel_futures=True)
提示
- SSH 批量命令:每连接一个线程,注意
MaxSessions与堡垒机限制 - 写操作(重启服务、删镜像)降低并发,必要时串行
- 结果写入 CSV/JSON 时,收集完再一次性落盘,避免多线程写同一文件
- 线程池无法可靠终止卡住的 Python 线程;SSH/HTTP/subprocess 必须各自设置连接、读取和执行超时
小结
| API | 用途 |
|---|---|
ThreadPoolExecutor(max_workers=N) | 限制并发 |
submit(fn, *args) | 提交单任务,返回 Future |
as_completed(futures) | 按完成顺序迭代 |
map(fn, iterable) | 批量映射,保序 |
下一章介绍 asyncio + httpx,在单线程内高并发处理大量 HTTP 请求。