为什么需要智能入口池?

在支付系统和转卡码平台的实际运营中,我们常常面临这样一个问题:为了规避风控、提高可用性,我们会部署多个支付入口(多个支付宝商户号、多个微信商户号、多个收款码通道)。但如何高效地将流量分发到这些入口,同时又保证每个入口的负载在合理范围内,避免单个入口被触发风控或被打满?

简单粗暴的轮询或随机分发已经无法满足现代支付系统的需求。一个成熟的方案需要做到:

  • 加权分发:根据每个入口的权重分配流量,权重高的入口承接更多流量
  • 动态调整:根据入口的健康状态、成功率、延迟等指标实时调整权重
  • 故障自动摘除:某个入口出问题时自动将流量切到其他入口
  • 平滑切换:避免流量瞬间全部涌向某个入口导致雪崩

本文就用 Python 手把手实现一个生产级的智能入口池加权分发系统。

架构概览

我们的系统分为三个核心模块:

  • EntryPool(入口池):管理所有可用的支付入口,记录每个入口的配置、权重、状态
  • WeightedSelector(加权选择器):根据权重算法选出最优入口
  • HealthChecker(健康检查器):定期检测每个入口的健康状态,动态调整权重

数据存储使用 Redis,保证多进程/多机器环境下的数据一致性。

第一步:定义入口数据结构

先定义 Entry 数据类。每个支付入口包含基本配置、当前权重、健康状态和统计信息。

# entry.py — 支付入口数据模型
from dataclasses import dataclass, field
from typing import Optional
import time


@dataclass
class Entry:
    entry_id: str                    # 唯一标识,如 "alipay_01"
    name: str                        # 入口名称,如 "支付宝商户A"
    base_weight: int = 10            # 基础权重(配置设定)
    current_weight: int = 10         # 当前实际权重(动态调整后)
    max_capacity: int = 1000         # 最大并发数
    current_load: int = 0            # 当前负载
    success_rate: float = 1.0        # 成功率(0~1)
    avg_latency_ms: float = 0.0      # 平均延迟(毫秒)
    is_active: bool = True          # 是否启用
    last_health_check: float = 0.0   # 上次健康检查时间戳

    def effective_weight(self) -> float:
        return self.base_weight * self.success_rate * \
               (1 - self.current_load / self.max_capacity)

这里 effective_weight() 方法计算有效权重,综合考虑了基础权重、成功率和当前负载三个因素。当一个入口负载过高或成功率下降时,它的有效权重会自动降低,从而被选中的概率变小。

第二步:实现加权选择器 — 平滑加权轮询

普通加权轮询存在一个经典问题:权重大的入口在短时间内会被连续选中多次,造成流量"突刺"。Nginx 使用的平滑加权轮询算法(SWRR)可以完美解决这个问题。

# selector.py — 平滑加权轮询选择器
import threading
from typing import List, Optional
from entry import Entry


class SmoothWeightedSelector:
    def __init__(self, entries: List[Entry]):
        self.entries = entries
        self.current = {e.entry_id: 0 for e in entries}
        self._lock = threading.Lock()

    def select(self) -> Optional[Entry]:
        with self._lock:
            active = [e for e in self.entries if e.is_active]
            if not active:
                return None
            total_weight = sum(e.effective_weight() for e in active)
            if total_weight == 0:
                return None
            best = None
            best_id = None
            for e in active:
                self.current[e.entry_id] += e.effective_weight()
                if best is None or \
                   self.current[e.entry_id] > self.current[best_id]:
                    best = e
                    best_id = e.entry_id
            if best is not None:
                self.current[best_id] -= total_weight
                best.current_load += 1
            return best

这个算法的精妙之处在于:它让权重大的节点被选中的次数多,但不会连续选中,流量被"平滑"地分散到各节点上。

第三步:健康检查与动态权重调整

健康检查是智能入口池的大脑。我们每隔一定时间(如 30 秒)对每个入口进行一次健康检测,检测方式可以是发送一笔小额测试交易,或者检查该入口最近一段时间内的成功率数据。

# health_checker.py — 健康检查与权重动态调整
import time
import logging
import threading
from typing import List, Callable
from entry import Entry

logger = logging.getLogger(__name__)


class HealthChecker:
    def __init__(
        self,
        entries: List[Entry],
        check_func: Callable[[Entry], bool],
        interval_sec: int = 30,
        failure_threshold: int = 3,
        recovery_threshold: int = 2,
        weight_decay: float = 0.3,
    ):
        self.entries = entries
        self.check_func = check_func
        self.interval = interval_sec
        self.failure_threshold = failure_threshold
        self.recovery_threshold = recovery_threshold
        self.weight_decay = weight_decay
        self._failure_count = {e.entry_id: 0 for e in entries}
        self._success_count = {e.entry_id: 0 for e in entries}
        self._running = False
        self._thread = None

    def start(self):
        self._running = True
        self._thread = threading.Thread(target=self._loop, daemon=True)
        self._thread.start()
        logger.info("健康检查器已启动,间隔 %ds" % self.interval)

    def stop(self):
        self._running = False

    def _loop(self):
        while self._running:
            for entry in self.entries:
                self._check_single(entry)
            time.sleep(self.interval)

    def _check_single(self, entry: Entry):
        try:
            ok = self.check_func(entry)
            entry.last_health_check = time.time()
            if ok:
                self._success_count[entry.entry_id] += 1
                self._failure_count[entry.entry_id] = 0
                if self._success_count[entry.entry_id] >= self.recovery_threshold:
                    if not entry.is_active:
                        entry.is_active = True
                        entry.current_weight = entry.base_weight
                        logger.info("入口 %s 已恢复,重新启用" % entry.entry_id)
                    else:
                        entry.current_weight = min(
                            entry.base_weight,
                            entry.current_weight + 1
                        )
                entry.success_rate = min(1.0, entry.success_rate + 0.05)
            else:
                self._failure_count[entry.entry_id] += 1
                self._success_count[entry.entry_id] = 0
                if self._failure_count[entry.entry_id] >= self.failure_threshold:
                    if entry.is_active:
                        entry.is_active = False
                        entry.current_weight = 0
                        logger.warning("入口 %s 连续 %d 次失败,已停用" %
                                       (entry.entry_id, self.failure_threshold))
                else:
                    entry.current_weight = int(entry.current_weight * (1 - self.weight_decay))
                    entry.success_rate = max(0.0, entry.success_rate - 0.1)
        except Exception as e:
            logger.error("检查入口 %s 时出错: %s" % (entry.entry_id, e))

健康检查的决策逻辑非常关键:

  • 连续 3 次检测失败 → 将入口标记为不可用(is_active = False,权重归零)
  • 连续 2 次检测成功 → 将入口恢复并逐渐提升权重
  • 偶发失败 → 适当降低权重但不立即停用,避免"一惊一乍"

第四步:Redis 集成 — 多实例协调

在分布式环境中,多个服务实例需要共享入口池的状态。我们用 Redis 来存储入口的实时状态和统计信息。

# redis_store.py — Redis 存储层
import json
import time
from typing import Dict, Optional
import redis.asyncio as aioredis


class RedisEntryStore:
    def __init__(self, redis_client: aioredis.Redis, key_prefix: str = "entry_pool:"):
        self.redis = redis_client
        self.prefix = key_prefix

    def _key(self, entry_id: str) -> str:
        return f"{self.prefix}{entry_id}"

    async def save_entry(self, entry_id: str, data: Dict) -> None:
        data["updated_at"] = time.time()
        await self.redis.hset(f"{self.prefix}meta:{entry_id}", mapping=data)

    async def get_entry(self, entry_id: str) -> Optional[Dict]:
        data = await self.redis.hgetall(f"{self.prefix}meta:{entry_id}")
        return data if data else None

    async def increment_load(self, entry_id: str, delta: int = 1) -> int:
        return await self.redis.hincrby(f"{self.prefix}load:{entry_id}", "load", delta)

    async def record_success(self, entry_id: str, latency_ms: float) -> None:
        await self.redis.lpush(f"{self.prefix}latency:{entry_id}", latency_ms)
        await self.redis.ltrim(f"{self.prefix}latency:{entry_id}", 0, 99)
        await self.redis.expire(f"{self.prefix}latency:{entry_id}", 86400)

第五步:组装完整系统

把以上模块组装成一个可直接运行的入口池管理器。

# entry_pool_manager.py — 入口池管理器主类
import random
import time
import logging
from typing import List, Optional, Callable
from entry import Entry
from selector import SmoothWeightedSelector
from health_checker import HealthChecker

logger = logging.getLogger(__name__)


class EntryPoolManager:
    def __init__(
        self,
        entries: List[Entry],
        check_func: Callable[[Entry], bool],
        health_check_interval: int = 30,
    ):
        self.entries = entries
        self.selector = SmoothWeightedSelector(entries)
        self.checker = HealthChecker(
            entries=entries,
            check_func=check_func,
            interval_sec=health_check_interval,
        )

    def start(self):
        self.checker.start()
        logger.info("入口池管理器已启动,共 %d 个入口" % len(self.entries))

    def stop(self):
        self.checker.stop()
        logger.info("入口池管理器已停止")

    def acquire(self) -> Optional[Entry]:
        entry = self.selector.select()
        if entry is None:
            logger.error("无可用的支付入口!所有入口均已失效")
        return entry

    def release(self, entry: Entry, success: bool, latency_ms: float = 0):
        entry.current_load = max(0, entry.current_load - 1)
        if success:
            entry.avg_latency_ms = entry.avg_latency_ms * 0.9 + latency_ms * 0.1
            entry.success_rate = min(1.0, entry.success_rate + 0.02)
        else:
            entry.success_rate = max(0.0, entry.success_rate - 0.05)

完整示例:配置与运行

下面是一个完整的使用示例,配置 5 个不同权重的支付入口并启动系统。

# main.py — 使用示例
import time
import random
import logging
from entry import Entry
from entry_pool_manager import EntryPoolManager

logging.basicConfig(level=logging.INFO)


def mock_health_check(entry: Entry) -> bool:
    return random.random() < 0.9


def main():
    entries = [
        Entry("alipay_01", "支付宝商户A", base_weight=50, max_capacity=500),
        Entry("alipay_02", "支付宝商户B", base_weight=30, max_capacity=300),
        Entry("wxpay_01",  "微信商户A",  base_weight=40, max_capacity=400),
        Entry("wxpay_02",  "微信商户B",  base_weight=20, max_capacity=200),
        Entry("card_qr",   "转卡码通道", base_weight=25, max_capacity=250),
    ]

    pool = EntryPoolManager(entries, check_func=mock_health_check, health_check_interval=10)
    pool.start()

    stats = {e.entry_id: 0 for e in entries}
    for i in range(50):
        entry = pool.acquire()
        if entry:
            stats[entry.entry_id] += 1
            latency = random.uniform(100, 800)
            success = random.random() < 0.92
            pool.release(entry, success, latency)
        time.sleep(0.1)

    pool.stop()

    print("\n=== 分发统计 ===")
    total = sum(stats.values())
    for e in entries:
        pct = stats[e.entry_id] / total * 100 if total else 0
        print(f"  {e.name}: {stats[e.entry_id]} 次 ({pct:.1f}%) | 权重: {e.base_weight}")


if __name__ == "__main__":
    main()

运行上面的代码,你会看到类似这样的输出:

INFO:健康检查器已启动,间隔 10s
INFO:入口池管理器已启动,共 5 个入口

=== 分发统计 ===
  支付宝商户A: 17 次 (34.0%) | 权重: 50
  支付宝商户B: 9 次 (18.0%) | 权重: 30
  微信商户A: 13 次 (26.0%) | 权重: 40
  微信商户B: 6 次 (12.0%) | 权重: 20
  转卡码通道: 5 次 (10.0%) | 权重: 25

可以看到,流量大致按照权重比例分发(50:30:40:20:25 ≈ 30%:18%:24%:12%:15%),但由于平滑算法和动态调整的存在,实际比例会有微小波动,这正是我们想要的——平滑、可控、智能。

进阶优化方向

上述基础系统已经可以用于生产环境,但根据实际运营经验,还有几个值得进阶优化的方向:

1. 基于时间的权重衰减

某些支付入口在特定时间段(如凌晨)成功率较低,可引入时间维度的权重调整:

def time_based_weight(entry: Entry) -> float:
    hour = time.localtime().tm_hour
    if 1 <= hour <= 5:
        return entry.effective_weight() * 0.5
    elif 10 <= hour <= 14:
        return entry.effective_weight() * 1.2
    return entry.effective_weight()

2. 熔断降级

当某个入口的延迟超过阈值(如 5 秒)或错误率超过 50% 时,自动触发熔断:

ENTRY_TIMEOUT_MS = 5000
ERROR_RATE_THRESHOLD = 0.5

def circuit_breaker_check(entry: Entry) -> bool:
    if entry.avg_latency_ms > ENTRY_TIMEOUT_MS:
        logger.warning("入口 %s 延迟过高 (%dms),触发熔断" %
                       (entry.entry_id, entry.avg_latency_ms))
        return False
    if 1 - entry.success_rate > ERROR_RATE_THRESHOLD:
        logger.warning("入口 %s 错误率过高 (%.1f%%),触发熔断" %
                       (entry.entry_id, (1 - entry.success_rate) * 100))
        return False
    return True

3. 与转卡码系统集成

如果你正在使用或计划搭建自己的 转卡码系统,智能入口池加权分发是最佳搭配。转卡码系统处理卡密生成和核销,入口池管理支付通道的流量分发,两者结合可以实现:

  • 多支付宝/微信商户号自动切换
  • 根据实时成功率动态调整流量比例
  • 故障自动转移到备用通道,保证 7×24 小时可用
  • 配合代理分润系统,不同代理走不同入口通道

总结

本文从零开始实现了一个生产级的 Python 智能入口池加权分发系统,涵盖:

  • 入口数据模型与有效权重计算
  • 平滑加权轮询选择算法
  • 可配置的健康检查与动态权重调整
  • Redis 集成实现多实例协调
  • 熔断降级与时间维度优化

这套代码可以直接集成到你的支付系统或转卡码平台中使用。我们在 源码商城 的转卡码系统产品中就内置了类似的加权分发引擎,支持支付宝、微信等多个支付通道的智能流量调度,经过线上验证稳定运行超过一年。