影刀RPA店群自动化稳定性实战:流程异常处理与自愈系统设计

影刀RPA店群自动化稳定性实战:流程异常处理与自愈系统设计

自动化系统上线后的前两个月,我养成了一种条件反射:
半夜醒来摸手机看告警群,看到“失败率飙升”就浑身一紧。

不是流程写错了。
而是线上环境的复杂程度,远远超出测试时的想象。

网络闪断、页面改版、登录态过期、平台限流、验证码突然弹出——
这些事任何一件单独发生,都不算什么大问题。
但当几十个店铺、三个平台、几百个任务同时跑的时候,这些小概率事件的叠加,会让系统变得千疮百孔。

这篇文章不写怎么让流程不报错。
我要复盘的是,报错之后,系统怎么扛过去。


一、异常不是敌人,措手不及才是

刚开始做自动化时,我们对异常的认知很粗糙:
流程返回失败,就在 Python 里打一条日志,然后人工上去看看怎么回事。

店铺少的时候,一天也就三五次失败,还能应付。
店铺数上来以后,每天上百次异常,种类五花八门。
有次一个拼多多店铺因为图片格式不对,连续失败了 40 次,重试了 40 次——
代理 IP 的流量被耗光,才被人发现。

从那之后我们花了很大力气重新梳理异常。
不是去消灭异常,而是让系统在异常面前有章法。


二、异常分类:每种错误都必须知道“接下来该怎么办”

我们把线上遇到的异常,归纳成四个大类:

picture.image

picture.image

  1. 瞬时故障:网络超时、连接重置、服务暂时不可用。特点是重试几次通常能恢复。
  2. 环境故障:浏览器崩溃、内存不足、进程僵死。需要回收资源后重试。
  3. 平台主动拦截:登录态失效、验证码、限流、风控验证。需要恢复环境后重试,且可能需要更换 IP 或切换账号状态。
  4. 业务错误:商品数据不合法、店铺权限不足、接口参数错误。重试无意义,必须停止并通知运营。

picture.image 这个分类体系落地成了一段配置和一套 Python 枚举:

from enum import Enum

class ExceptionCategory(Enum):
    TRANSIENT = "transient"       # 可自动重试
    ENVIRONMENT = "environment"   # 需回收资源后重试
    
![picture.image](https://p3-volc-community-sign.byteimg.com/tos-cn-i-tlddhu82om/94f549cbba3d413cb8aa7a278c963fa5~tplv-tlddhu82om-image.image?=&rk3s=8031ce6d&x-expires=1785835614&x-signature=ddu8RzX8Vn0pzL6RO7rrZbPHQI0%3D)
    PLATFORM_GATE = "platform"    # 需状态恢复后重试
    BUSINESS = "business"         # 不可重试,需人工介入

picture.image 接下来要做的,就是让每个影刀流程的失败回调,带上一个明确的错误码。
错误码到类别的映射,维护在一个独立的 YAML 文件里,方便随时调整。

picture.image

error_mapping:
  pdd:
    "E_TIMEOUT": transient
    "E_CONN_RESET": transient
    "E_BROWSER_CRASH": environment
    "E_LOGIN_EXPIRED": platform
    "E_CAPTCHA": platform
    "E_RATE_LIMIT": platform
    "E_INVALID_PARAM": business
  temu:
    # 类似结构

这样 Python 侧拿到错误码后,就知道该走重试、走资源回收、还是直接告警。


三、重试策略的核心不是次数,是间隔和判断

很多人一说重试,就定“重试 3 次”。
至于间隔多少、什么条件下停,一概没想。

我们踩过的坑之一,就是“连续重试导致平台限流升级”。
一个拼多多店铺因网络超时触发了重试,间隔 2 秒就重试一次,连续 3 次不仅没恢复,反而被平台判定为高频异常请求,封了 IP 一小时。

后来我们把重试策略改成了指数退避 + 最大延迟上限 + 错误类型感知

import asyncio
import random

class RetryPolicy:
    def __init__(self, max_retries=3, base_delay=10, max_delay=120):
        self.max_retries = max_retries
        self.base_delay = base_delay
        self.max_delay = max_delay

    def delay(self, attempt: int) -> float:
        # 指数退避,加入随机抖动
        exp = min(self.base_delay * (2 ** attempt), self.max_delay)
        return exp * (0.5 + random.random())

在重试前,还会根据异常类别做前置检查:

  • 如果是 ENVIRONMENT 类异常,必须先回收浏览器实例,再进入重试。
  • 如果是 PLATFORM_GATE,且错误码是 E_LOGIN_EXPIRED,重试前必须触发一个“重新登录”的恢复流程,否则重试毫无意义。
  • BUSINESS 类异常直接跳过重试,标记失败并通知。

重试的每一次调用,都要生成新的 attempt_id,并在日志里清晰记录每次重试的原因和耗时。
这样在复盘时,可以清楚地看到某个任务在重试中是否在好转还是在恶化。


四、自愈的第一步:环境恢复

有很多错误,不是重跑流程就能解决的。
浏览器已经崩了,登录态已经掉了,IP 已经被限制了。

这时候需要的不是重试,而是一套环境恢复动作。

我们设计了一个 EnvRecovery 模块,它会根据异常类别,执行对应的恢复操作。

class EnvRecovery:
    def __init__(self, browser_pool, session_mgr, proxy_mgr):
        self.browser_pool = browser_pool
        self.session_mgr = session_mgr
        self.proxy_mgr = proxy_mgr

    async def recover(self, shop_id: str, category: ExceptionCategory, error_code: str):
        if category == ExceptionCategory.ENVIRONMENT:
            # 强制销毁该店铺的浏览器实例,标记为下次需重建
            await self.browser_pool.destroy(shop_id)
            return True
        elif category == ExceptionCategory.PLATFORM_GATE:
            if error_code == "E_LOGIN_EXPIRED":
                await self.session_mgr.invalidate(shop_id)
                # 触发一个登录恢复任务
                return "REQUIRE_LOGIN"
            elif error_code == "E_RATE_LIMIT":
                await self.proxy_mgr.rotate(shop_id)
                return True
        return False

恢复动作本身是异步的、幂等的。
同一个店铺如果连续 3 次恢复失败,就会触发 P1 告警,避免陷入无限恢复循环。

这里有一个真实的教训:
早期我们没有对恢复次数做上限,一个 TEMU 店铺因为登录失败,反复触发重新登录,结果连续登录失败触发了平台的安全锁定,店铺被冻结了三天。
自那以后,所有自愈动作都带上了熔断计数器。

class CircuitBreaker:
    def __init__(self, threshold=3, cooldown_seconds=600):
        self.failures = {}
        self.threshold = threshold
        self.cooldown = cooldown_seconds

    def record_failure(self, shop_id: str):
        now = time.time()
        if shop_id not in self.failures:
            self.failures[shop_id] = []
        self.failures[shop_id].append(now)
        # 只保留冷却窗口内的失败记录
        self.failures[shop_id] = [t for t in self.failures[shop_id] if now - t < self.cooldown]

    def is_open(self, shop_id: str) -> bool:
        return len(self.failures.get(shop_id, [])) >= self.threshold

熔断器打开后,该店铺的所有新任务都会被暂停,直接进入人工处理队列,不再自动重试。


五、断点续跑:让我们又爱又恨的能力

断点续跑,意思是一个长流程从失败的地方接着跑,而不是从头开始。
这个需求在运营侧呼声极高——一个商品上架流程跑了 15 分钟,在最后一步提交时失败,从头跑就是巨大的时间浪费。

我们在技术上论证了很久,最终决定:只对部分可控的步骤做断点续跑,不追求全流程覆盖

原因是,浏览器环境的状态是不可信的。
你以为流程断在了第三步,但实际可能第二步的某个副作用已经生效了,只是没拿到反馈。
从断点重入,风险太高。

我们采用的折中方案是:
把一个大流程拆成多个“不可分割的原子步骤组”,每个组之间设置断点。
断点处,流程必须把当前状态(已完成的步骤、已获取的数据)通过回调写入外部存储。
如果流程中断,Python 侧可以根据上次保存的状态,决定从哪个原子组重新开始,而不是从零起步。

class FlowCheckpoint:
    def __init__(self, redis_client):
        self.redis = redis_client

    async def save(self, task_id: str, step: str, context: dict):
        key = f"flow:checkpoint:{task_id}"
        await self.redis.hset(key, step, json.dumps(context))
        await self.redis.expire(key, 3600)

    async def load_latest(self, task_id: str) -> dict:
        key = f"flow:checkpoint:{task_id}"
        data = await self.redis.hgetall(key)
        if not data:
            return {"step": "start", "context": {}}
        latest_step = max(data.keys())
        return {"step": latest_step, "context": json.loads(data[latest_step])}

一个实际的例子:拼多多商品上架,我们拆成“登录验证 → 填写基本信息 → 上传图片 → 填写规格库存 → 提交”五个原子组。
如果在“填写规格库存”阶段失败,且前面“上传图片”的断点已保存,那重跑时直接从“填写规格库存”开始,省掉前面 10 分钟。
但如果图片上传内部失败,断点就退回到“填写基本信息”,保证不出现“有规格但没图片”的脏数据。

这个设计谈不上完美,但在我们的场景里,把平均重跑时间从 8 分钟压到了 2 分钟以内,运维压力明显降低。


六、超时与看门狗:别让一个任务拖死整个节点

这个问题其实在任务调度那篇里提到过,但从异常处理的角度看,有必要再展开。

一个影刀流程可能因为任何原因卡住——
弹出了意料之外的对话框、页面无限加载、JS 死循环。

如果执行节点不加控制,一个任务可能占用浏览器槽位数小时。
最终导致节点上所有槽位被僵尸任务占满,其他任务全部饥饿。

我们做了两层超时控制:

  1. 业务超时:每个动作在配置文件里预设了最大执行时间,比如上传商品 300 秒,改价 120 秒。
  2. 系统硬超时:无论什么动作,最大执行不超过 600 秒。超过后直接 SIGKILL 影刀进程对应的浏览器渲染进程。
class HardTimeoutGuard:
    def __init__(self, timeout: int, on_timeout):
        self.timeout = timeout
        self.on_timeout = on_timeout

    async def __aenter__(self):
        self.task = asyncio.create_task(self._watch())
        return self

    async def __aexit__(self, exc_type, exc, tb):
        self.task.cancel()

    async def _watch(self):
        await asyncio.sleep(self.timeout)
        await self.on_timeout()

on_timeout 回调里,不仅停止影刀流程,还要回收浏览器、记录超时日志、并将任务标记为 TIMEOUT 而非 FAILED
TIMEOUT 状态的任务,在分析平台上会单独归类,方便评估各平台的动作耗时是否存在恶化趋势。


七、自愈系统的闭环:从发现到恢复,再到学习

自愈不是简单的“失败-重试-恢复”。
真正有价值的是,恢复之后,系统能不能记住教训。

我们做了一个简单的异常知识库:
每当一种新的错误码出现,运维人员会在处理完后,把对应的恢复策略更新到配置里。
但更重要的是,系统会统计每种异常的发生频率和恢复成功率。

如果某种异常的恢复成功率持续低于 50%,比如某平台的登录失效恢复总是不成功,
调度器就会逐步降低该平台的任务权重,直到问题被人工根治。

class PlatformHealthScorer:
    def __init__(self, db_session):
        self.db = db_session

    async def get_score(self, platform: str) -> float:
        # 查询最近1小时的恢复成功率
        stats = await self.db.query(
            "SELECT SUM(CASE WHEN recovered=1 THEN 1 ELSE 0 END) / COUNT(*) "
            "FROM recovery_log WHERE platform=:platform AND created_at > :since",
            {"platform": platform, "since": time.time() - 3600}
        )
        return float(stats or 1.0)

    async def adjust_capacity(self, platform: str):
        score = await self.get_score(platform)
        if score < 0.5:
            # 将任务分发权重降低50%
            await self.redis.set(f"rpa:weight:{platform}", 0.5)
        else:
            await self.redis.delete(f"rpa:weight:{platform}")

调度器在分发任务时,会读取这个权重,减少向“不健康”平台的投喂。
这个机制在某次 TEMU 大规模改版导致大量流程失效的时候,避免了任务的无效堆积,也保护了代理 IP 不被频繁标记。


八、异常处理的代码组织结构

为了避免异常处理逻辑散落在各处,我们把所有相关逻辑收敛到了 FlowExecutor 这个核心类里。
它对外只暴露一个 execute(task) 方法,内部串联了重试、恢复、熔断、超时等全部逻辑。

这是大概的结构,简化了很多业务细节:

class FlowExecutor:
    def __init__(self, retry_policy, recovery, breaker, timeout, checkpoint):
        self.retry_policy = retry_policy
        self.recovery = recovery
        self.breaker = breaker
        self.timeout = timeout
        self.checkpoint = checkpoint

    async def execute(self, task: Task) -> TaskResult:
        if self.breaker.is_open(task.shop_id):
            return TaskResult.aborted("circuit_breaker_open")

        attempt = 0
        while attempt <= self.retry_policy.max_retries:
            try:
                # 从最近的断点获取上下文
                resume = await self.checkpoint.load_latest(task.id)
                async with HardTimeoutGuard(self.timeout, self._on_timeout):
                    result = await self._run_flow(task, resume)
                if result.success:
                    return TaskResult.completed(result.data)
                else:
                    raise FlowException(result.error_code, result.error_msg)
            except FlowException as e:
                category = classify(e.error_code, task.platform)
                if category == ExceptionCategory.BUSINESS:
                    return TaskResult.failed(e)
                # 尝试恢复环境
                await self.recovery.recover(task.shop_id, category, e.error_code)
                self.breaker.record_failure(task.shop_id)
                if self.breaker.is_open(task.shop_id):
                    return TaskResult.aborted("circuit_breaker_open")
                attempt += 1
                delay = self.retry_policy.delay(attempt)
                await asyncio.sleep(delay)
            except Exception as e:
                # 未知异常,快速失败并告警
                await self._alert_unknown(task, e)
                return TaskResult.error(e)

        return TaskResult.max_retries_exceeded()

真实代码比这个多了很多监控埋点和日志打印,但骨架就是这样。
把异常处理的复杂度关在一个类里,让调度器和守护进程不用关心细节,这是我们花了不少时间才重构出来的清晰边界。


九、不止于恢复,而是让系统更抗造

自愈系统上线后,我们并没有高枕无忧。
每个月的复盘里,我们还是会挑出那些“本该自愈但没自愈”的案例来分析。

有的原因是错误码归类错了——
TEMU 某个接口返回的 -10001 表面像网络超时,实际是商品类目失效,属于业务错误。
我们把分类一改,重试就不再触发,任务直接转人工,反而处理更快。

有的是自愈动作太“暴力”——
一遇到登录过期就销毁整个浏览器实例重建,影响同机器上其他正在运行的店铺。
后来改成只清理登录态的 Session 文件,不销毁浏览器,粒度更细。

自动化的稳定性不是一蹴而就的。
它是一个持续观察、持续调优的过程,靠的是数据和经验,不是一次性的代码发布。

很多团队在最开始都会忽略异常体系的建设,把所有精力放在正常流程上。
直到系统被层出不穷的边缘情况打到千疮百孔,才回头补这个课。
如果让我重新来一次,我会在第三个店铺上线前,就把重试、恢复、熔断这套架子搭好。


十、写在最后

做店群自动化这几年,我最大的体会是:
能让流程跑起来的代码,只占整个系统工作量的 30%。
剩下的 70%,全是让系统在出事时还能活着、还能恢复、还能给运维留条活路的工作。

异常处理和自愈,就是这 70% 里最核心的一环。

它不酷,不动辄“智能化”,甚至大部分时间都是默默无闻的。
但它决定了你的系统,是一碰就碎的玩具,还是一个能扛住真实世界千锤百炼的生产工具。

如果你正好在为线上各种莫名奇妙的失败头疼,希望这篇复盘能给你一些切实的启发。

作者:林焱

0
0
0
0
评论
未登录
暂无评论