您的位置:首页 >Pythonasyncio异步并发与多固定出口IP调度实战
发布于2026-08-13 阅读(0)
扫一扫,手机访问
单事件循环下管理多个出口,最大的挑战是状态一致性:多个协程同时 acquire/release 出口,如果用普通的 dict 管理分配状态,协程切换时机不对会导致同一个出口被超量分配。必须用 asyncio.Lock 或 asyncio.Semaphore 做并发保护。

第二个挑战是健康检查不阻塞:同步场景健康检查在后台线程跑,阻塞不影响主流程;异步场景下如果健康检查协程卡在 await 上,会拖慢整个事件循环的调度。需要给健康检查加超时和熔断。
aiohttp 支持在应用层配置 HTTP 转发地址,但这有两个问题:一是每条请求都经过中间节点转发,多一跳延迟;二是某些对端服务会检测转发协议头(如 Via、X-Forwarded-For),判定为非直连。
更稳妥的方式是源地址绑定:在 TCP 层将连接绑定到指定的本地 IP 地址,对端看到的是该 IP 的直连请求,没有转发特征。
import asyncio
import aiohttp
from dataclasses import dataclass, field
from typing import Optional
import time
@dataclass
class ExitConfig:
"""单个固定出口的配置"""
name: str # 出口标识
source_ip: str # 绑定的本地源 IP
max_concurrent: int = 100 # 该出口最大并发连接数
enabled: bool = True
@dataclass
class ExitState:
"""出口运行时状态"""
config: ExitConfig
connector: Optional[aiohttp.TCPConnector] = None
healthy: bool = True
active_count: int = 0 # 当前活跃连接数
last_check: float = 0.0 # 上次健康检查时间
fail_count: int = 0 # 连续失败次数
total_requests: int = 0 # 总请求数
total_errors: int = 0 # 总错误数
@property
def error_rate(self) -> float:
if self.total_requests == 0:
return 0.0
return self.total_errors / self.total_requests
@property
def a vailable(self) -> bool:
"""是否可分配:健康 + 未达并发上限"""
return (
self.config.enabled
and self.healthy
and self.active_count < self.config.max_concurrent
)async def create_connector(exit_config: ExitConfig) -> aiohttp.TCPConnector:
"""为指定出口创建绑定源 IP 的 TCPConnector"""
connector = aiohttp.TCPConnector(
local_addr=(exit_config.source_ip, 0), # 绑定源 IP,端口随机
limit=exit_config.max_concurrent, # 连接池上限
limit_per_host=50, # 单对端最大连接
keepalive_timeout=30, # keepalive 超时
enable_cleanup_closed=True, # 关闭的连接及时清理
ttl_dns_cache=300, # DNS 缓存 5 分钟
)
return connectorlocal_addr 参数是关键——它在 socket 层调用 bind() 绑定到指定本地 IP,所有通过该 Connector 发出的连接都从这个 IP 出去。用 ss -tnp 可以验证:
# 查看指定源 IP 的活跃连接
ss -tnp | grep "10.0.1.100"
# 输出示例:
# ESTAB 0 0 10.0.1.100:54321 93.184.216.34:443class AsyncExitPool:
"""异步多固定出口 IP 调度池"""
def __init__(self, exit_configs: list[ExitConfig]):
self._exits: dict[str, ExitState] = {}
self._lock = asyncio.Lock()
self._init_exits(exit_configs)
def _init_exits(self, configs: list[ExitConfig]):
for cfg in configs:
self._exits[cfg.name] = ExitState(config=cfg)
async def start(self):
"""初始化所有出口的 Connector"""
for state in self._exits.values():
state.connector = await create_connector(state.config)
async def acquire(self, preferred: Optional[str] = None) -> tuple[str, aiohttp.TCPConnector]:
"""
获取一个可用出口的 Connector
preferred: 优先使用的出口名称(None 表示自动选择负载最低的)
返回: (出口名称, Connector)
"""
async with self._lock:
# 优先尝试指定出口
if preferred and preferred in self._exits:
state = self._exits[preferred]
if state.a vailable:
state.active_count += 1
return preferred, state.connector
# 自动选择:负载最低的健康出口
candidates = [s for s in self._exits.values() if s.a vailable]
if not candidates:
raise RuntimeError("无可用出口")
# 按 active_count / max_concurrent 比率排序,选负载最低的
candidates.sort(key=lambda s: s.active_count / s.config.max_concurrent)
chosen = candidates[0]
chosen.active_count += 1
return chosen.config.name, chosen.connector
async def release(self, name: str, success: bool = True):
"""归还出口"""
async with self._lock:
state = self._exits.get(name)
if state:
state.active_count = max(0, state.active_count - 1)
state.total_requests += 1
if not success:
state.total_errors += 1
state.fail_count += 1
if state.fail_count >= 3:
state.healthy = False
print(f"[WARN] 出口 {name} 连续失败 {state.fail_count} 次,标记为不健康")
else:
state.fail_count = 0 # 成功则重置失败计数
async def get_stats(self) -> list[dict]:
"""获取所有出口状态"""
async with self._lock:
return [
{
"name": s.config.name,
"source_ip": s.config.source_ip,
"healthy": s.healthy,
"active": s.active_count,
"max": s.config.max_concurrent,
"usage": f"{s.active_count}/{s.config.max_concurrent}",
"total_req": s.total_requests,
"error_rate": f"{s.error_rate:.2%}",
"fail_count": s.fail_count,
}
for s in self._exits.values()
]async def fetch_with_exit(pool: AsyncExitPool, url: str, preferred: str = None):
"""使用固定出口发起异步请求"""
exit_name, connector = await pool.acquire(preferred)
success = False
try:
timeout = aiohttp.ClientTimeout(total=10, connect=5)
async with aiohttp.ClientSession(connector=connector, timeout=timeout) as session:
async with session.get(url) as resp:
if resp.status == 200:
success = True
return await resp.text()
else:
print(f"[WARN] {exit_name} -> {url} status={resp.status}")
return None
except Exception as e:
print(f"[ERROR] {exit_name} -> {url}: {e}")
return None
finally:
await pool.release(exit_name, success)
async def main():
configs = [
ExitConfig(name="us-01", source_ip="10.0.1.100", max_concurrent=100),
ExitConfig(name="us-02", source_ip="10.0.1.101", max_concurrent=100),
ExitConfig(name="eu-01", source_ip="10.0.2.100", max_concurrent=80),
]
pool = AsyncExitPool(configs)
await pool.start()
# 并发 500 个任务,自动分配到 3 个出口
urls = [f"https://httpbin.org/delay/{i % 3}" for i in range(500)]
tasks = [fetch_with_exit(pool, url) for url in urls]
results = await asyncio.gather(*tasks, return_exceptions=True)
# 查看出口使用情况
stats = await pool.get_stats()
for s in stats:
print(s)
asyncio.run(main())class ExitHealthChecker:
"""异步出口健康检查器"""
def __init__(self, pool: AsyncExitPool, check_url: str, interval: int = 30):
self.pool = pool
self.check_url = check_url # 健康检查用的对端地址
self.interval = interval # 检查间隔(秒)
self._task: Optional[asyncio.Task] = None
async def _check_one(self, name: str, state: ExitState):
"""检查单个出口"""
timeout = aiohttp.ClientTimeout(total=5, connect=3)
try:
async with aiohttp.ClientSession(connector=state.connector, timeout=timeout) as session:
start = time.monotonic()
async with session.get(self.check_url) as resp:
latency = (time.monotonic() - start) * 1000
if resp.status == 200:
async with self.pool._lock:
if not state.healthy:
print(f"[INFO] 出口 {name} 恢复健康 (latency={latency:.0f}ms)")
state.healthy = True
state.fail_count = 0
state.last_check = time.time()
return True
except Exception:
async with self.pool._lock:
state.fail_count += 1
if state.fail_count >= 3 and state.healthy:
state.healthy = False
print(f"[WARN] 出口 {name} 健康检查失败,标记不健康")
state.last_check = time.time()
return False
async def _run(self):
"""定时检查循环"""
while True:
tasks = [
self._check_one(name, state)
for name, state in self.pool._exits.items()
if state.config.enabled
]
await asyncio.gather(*tasks, return_exceptions=True)
await asyncio.sleep(self.interval)
def start(self):
"""启动健康检查后台协程"""
self._task = asyncio.create_task(self._run())
def stop(self):
if self._task:
self._task.cancel()当 acquire() 发现所有出口不可用时,有两种策略:
async def acquire_with_fallback(self, preferred: str = None) -> tuple[str, aiohttp.TCPConnector]:
"""带故障转移的获取出口"""
try:
return await asyncio.wait_for(self.acquire(preferred), timeout=3.0)
except (RuntimeError, asyncio.TimeoutError):
# 所有出口不可用,尝试重新激活被标记为不健康的出口
async with self._lock:
for state in self._exits.values():
if not state.healthy and state.config.enabled:
print(f"[WARN] 全部出口不可用,尝试重新激活 {state.config.name}")
state.healthy = True
state.fail_count = 0
# 重试一次
return await asyncio.wait_for(self.acquire(preferred), timeout=3.0)每个出口独立限速,避免单个 IP 请求过密触发对端限制:
class AsyncTokenBucket:
"""异步令牌桶:控制单个出口的请求速率"""
def __init__(self, rate: float, capacity: int):
self.rate = rate # 令牌生成速率(个/秒)
self.capacity = capacity # 桶容量
self._tokens = capacity # 当前令牌数
self._last_refill = time.monotonic()
self._lock = asyncio.Lock()
async def acquire(self, tokens: int = 1):
"""获取令牌,不足时等待"""
async with self._lock:
while self._tokens < tokens:
# 计算需要等待的时间
needed = tokens - self._tokens
wait_time = needed / self.rate
await asyncio.sleep(wait_time)
self._refill()
self._tokens -= tokens
def _refill(self):
now = time.monotonic()
elapsed = now - self._last_refill
self._tokens = min(self.capacity, self._tokens + elapsed * self.rate)
self._last_refill = nowclass ExitRotator:
"""出口轮换管理:按任务特征分配出口"""
def __init__(self, pool: AsyncExitPool):
self.pool = pool
self._buckets: dict[str, AsyncTokenBucket] = {}
def register_rate_limit(self, exit_name: str, rate: float, capacity: int):
"""为出口注册速率限制"""
self._buckets[exit_name] = AsyncTokenBucket(rate, capacity)
async def fetch(self, url: str, region: str = None) -> Optional[str]:
"""
按区域偏好选择出口并限速
region: 地区标识(如 us/eu),None 表示自动
"""
# 按地区筛选可用出口
preferred = f"{region}-01" if region else None
exit_name, connector = await self.pool.acquire(preferred)
# 获取该出口的令牌桶
bucket = self._buckets.get(exit_name)
if bucket:
await bucket.acquire() # 等待令牌
success = False
try:
timeout = aiohttp.ClientTimeout(total=10)
async with aiohttp.ClientSession(connector=connector, timeout=timeout) as session:
async with session.get(url) as resp:
if resp.status == 200:
success = True
return await resp.text()
self._tokens -= tokens
def _refill(self):
now = time.monotonic()
elapsed = now - self._last_refill
self._tokens = min(self.capacity, self._tokens + elapsed * self.rate)
self._last_refill = nowclass ExitRotator:
"""出口轮换管理:按任务特征分配出口"""
def __init__(self, pool: AsyncExitPool):
self.pool = pool
self._buckets: dict[str, AsyncTokenBucket] = {}
def register_rate_limit(self, exit_name: str, rate: float, capacity: int):
"""为出口注册速率限制"""
self._buckets[exit_name] = AsyncTokenBucket(rate, capacity)
async def fetch(self, url: str, region: str = None) -> Optional[str]:
"""
按区域偏好选择出口并限速
region: 地区标识(如 us/eu),None 表示自动
"""
# 按地区筛选可用出口
preferred = f"{region}-01" if region else None
exit_name, connector = await self.pool.acquire(preferred)
# 获取该出口的令牌桶
bucket = self._buckets.get(exit_name)
if bucket:
await bucket.acquire() # 等待令牌
success = False
try:
timeout = aiohttp.ClientTimeout(total=10)
async with aiohttp.ClientSession(connector=connector, timeout=timeout) as session:
async with session.get(url) as resp:
if resp.status == 200:
success = True
return await resp.text()
except Exception as e:
print(f"[ERROR] {exit_name}: {e}")
finally:
await self.pool.release(exit_name, success)
return Noneclass ExitMetrics:
"""出口池指标采集器"""
@staticmethod
def to_prometheus(stats: list[dict]) -> str:
lines = []
for s in stats:
labels = f'name="{s["name"]}",ip="{s["source_ip"]}"'
lines.append(f'exit_healthy{{{labels}}} {1 if s["healthy"] else 0}')
lines.append(f'exit_active_connections{{{labels}}} {s["active"]}')
lines.append(f'exit_max_connections{{{labels}}} {s["max"]}')
lines.append(f'exit_total_requests{{{labels}}} {s["total_req"]}')
# error_rate 需要解析
err = float(s["error_rate"].rstrip("%")) / 100
lines.append(f'exit_error_rate{{{labels}}} {err}')
return "n".join(lines) + "n"from aiohttp import web
async def metrics_handler(request: web.Request):
pool: AsyncExitPool = request.app["exit_pool"]
stats = await pool.get_stats()
text = ExitMetrics.to_prometheus(stats)
return web.Response(text=text, content_type="text/plain")
async def start_metrics_server(pool: AsyncExitPool, port: int = 9100):
app = web.Application()
app["exit_pool"] = pool
app.router.add_get("/metrics", metrics_handler)
runner = web.AppRunner(app)
await runner.setup()
site = web.TCPSite(runner, "0.0.0.0", port)
await site.start()
print(f"[INFO] 指标端点启动: http://0.0.0.0:{port}/metrics")误区一:异步场景直接复用同步的 ExitPool。 同步场景里,线程锁(threading.Lock)确实能把状态保护住;可一旦放到异步环境里,这套做法就不对劲了。线程锁会把整个事件循环卡住——某个协程一旦拿到锁,其他协程就没机会被切换上来,结果本来并发的异步流程,硬生生退化成了串行。正确的做法是换成 asyncio.Lock,让协程在 await 时自然让出事件循环。
误区二:每个请求创建新的 Connector。 Connector 内部维护连接池,频繁创建/销毁会导致 TCP 连接无法复用,每次都走完整握手。正确做法是每个出口一个 Connector,生命周期与 ExitPool 一致,只在关闭时清理。
误区三:健康检查不加超时。 健康检查协程如果 await 一个卡住的对端,会占用事件循环的调度槽。虽然 asyncio 不会完全阻塞(其他协程还能跑),但积压太多超时协程会导致事件循环响应变慢。每个健康检查请求必须设 total=5, connect=3 的超时。
误区四:用 Semaphore 替代 active_count 管理。asyncio.Semaphore 的作用,只是把总并发数卡住;它并不会告诉你,现在到底有多少资源正在被占用,哪一个出口又更“忙”。真正要做按负载分配,就得靠 active_count 这种显式计数方式,把每个出口的压力看得清清楚楚,进而优先选择负载最低的出口,而不是停留在简单轮询这一步。
# 1. 确认源地址绑定生效
ss -tnp | grep "10.0.1.100"
# 应看到大量 ESTAB 连接从该 IP 发出
# 2. 确认连接池复用(TIME_WAIT 少说明复用正常)
ss -s | grep "TIME-WAIT"
# 3. 查看出口使用均衡度
curl http://localhost:9100/metrics | grep exit_active
# 4. 确认事件循环没有阻塞(事件循环延迟 < 50ms)
# 在代码中插入:
# loop_start = time.monotonic()
# await asyncio.sleep(0)
# loop_delay = (time.monotonic() - loop_start) * 1000
# print(f"事件循环延迟: {loop_delay:.1f}ms")
# 5. 查看协程数量是否合理
# asyncio.all_tasks() 返回当前所有活跃协程
print(f"活跃协程数: {len(asyncio.all_tasks())}")异步场景下多固定出口 IP 管理的核心就三条:源地址绑定替代应用层转发(无转发特征、无额外跳数)、asyncio.Lock 保护状态一致性(协程级并发安全)、非阻塞健康检查与故障转移(不拖慢事件循环)。把 AsyncExitPool + 令牌桶限速 + Prometheus 指标这套搭起来,千级并发任务跨多个固定出口调度就能稳定跑。
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
正版软件
正版软件
正版软件
正版软件
正版软件
1
2
3
4
5
6
7
8
9