ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Serverless 异步处理模式实战:定时任务、消息队列与幂等重试全解析

Serverless 异步处理模式实战:定时任务、消息队列与幂等重试全解析 做了这么多年Serverless应用我逐渐摸清一个规律凡是能用异步解决的就别用同步硬扛。这话听起来像废话但真踩过坑的人才知道——在Serverless环境里超时、并发、成本这三个老大难问题几乎全部可以用异步处理模式化解。今天我想结合最近用serverless部署的一个定时任务给trae做每日自动签到来聊聊这块顺便把异步处理模式的套路、坑位和实操细节一次性说清楚。如果你正准备把业务搬到函数计算上或者已经在用Serverless但总被超时和重试搞得头疼这篇文章应该能帮你少走不少弯路。我会从最基础的“为什么必须异步”讲起再到消息队列选型、定时任务实现、幂等重试设计最后把常见坑位和排查手段都列出来。全程用实际项目说话不整虚的。1. 先搞清楚Serverless为什么绕不开异步1.1 一句话理解Serverless的“脾气”Serverless最核心的特征是“按需执行”和“有状态限制”。你的业务逻辑被打散成一个一个函数每个函数都像一个临时工来活了就启动一个执行环境干完活立刻销毁。这个设计带来了极致的弹性和按量计费但也带来三个天生痛点。第一执行时长受限。绝大多数云厂商的函数计算服务都有默认超时时间比如某云厂商的默认是3秒最大能调到15分钟甚至1小时取决于具体平台但超时时间越长计费越贵而且长任务会长时间占用执行环境导致并发配额被快速耗尽。第二执行环境有冷启动问题。函数长时间没有调用后下次触发需要重新初始化运行时这个冷启动时间可能从几百毫秒到几秒不等。第三函数默认是无状态的。你没法在内存里保存跨调用共享的数据也不该依赖本机文件系统做持久化。这些“脾气”决定了任何需要长时间处理、需要等待外部响应、需要跨多个服务协作的逻辑如果硬塞进一个同步请求里大概率要翻车。我见过太多人刚上手Serverless时把传统后端那套同步思维直接搬过来。比如一个下单接口同步调用一个函数函数里又去同步请求支付网关、同步写数据库、同步发短信通知。结果是接口响应慢得离谱用户等得骂娘函数超时导致订单状态不一致重试时又重复下单。这完全是自找的。1.2 同步请求的局限那些让人抓狂的“超时”咱们用最典型的场景来说——外部API调用。假设你的函数需要调用一个第三方服务来查询天气第三方服务平均响应时间200ms听起来挺快对不对但别忘了你还要加上网络传输时间、函数自身的业务处理时间、日志写入时间一旦第三方服务偶发变慢到2秒你的函数就悬了。在同步模式里客户端直接等待函数执行完毕。如果函数要跑10秒客户端就得干等10秒。更糟的是HTTP网关通常也有超时设置比如API网关默认60秒有些网关只有30秒。一旦函数执行超过网关超时网关会直接断开连接但函数可能还在后台执行——用户看到的是“请求失败”但实际上业务可能已经成功了这就是典型的“不确定状态”。我用源代码的教训换来的经验是任何可能超过3秒的操作都值得考虑拆成异步任务。这里说的“可能超过3秒”包括第三方API调用、复杂计算、批量数据处理、大文件生成、跨系统数据同步等。判断标准很简单用户真的需要立刻知道结果吗如果不需要那就异步化。异步处理模式的基本思路是把一个长任务拆成两个阶段第一阶段快速接收请求立刻返回“已受理”第二阶段在后台慢慢执行完成后通过回调、轮询或消息通知结果。Serverless天然适合做异步因为它的事件驱动模型和消息队列结合得非常自然。2. 异步处理的核心套路事件驱动与消息队列2.1 异步的本质把“干活”拆成“接单”和“执行”理解异步处理我想用一个生活化的类比你在餐厅点餐。同步模式像是你去柜台点单然后站在厨师旁边等菜厨师炒一道菜你端走一道。如果厨师做菜要20分钟你就得干站20分钟。异步模式则是你点完单拿一个号牌回座位玩手机。后厨接到订单做好之后叫号你再过去取餐。如果菜品太多后厨还会把订单按顺序排队分批处理。在这个类比里“接单”就是函数的第一阶段校验参数、生成任务ID、把任务写入队列然后立刻返回“我们收到订单了”。而“执行”阶段由另一个函数或者同一个函数的异步调用从队列里取出消息真正去干活。这两个阶段之间用消息队列作为缓冲。消息队列的作用不仅是解耦更关键的是削峰填谷。比如每天上午10点有一波签到高峰瞬时请求量可能是平时的20倍。如果你用同步函数直连数据库数据库大概率被冲垮。但如果先用队列把请求全部收下再让处理函数按照自己的速率从队列里拉取消息流量就像洪水被水库截住了下游压力变得可控。Serverless平台本身也原生支持异步调用。以常见的函数计算服务为例你可以配置“异步调用”模式客户端调用函数时平台立刻返回一个请求ID函数真正执行在后台进行。平台内部其实已经帮你接了一条队列。不过这种方式比较粗粒度只能满足“整个函数后台跑”的需求。如果你需要更精细的控制比如消息按优先级处理、失败后人工重放、多个任务并发协调那就得自己搭消息队列或事件总线。2.2 消息队列选型SQS、Kafka还是内存队列消息队列是异步处理的核心基础设施。选型上我把它分为三大类云托管队列、自建消息系统、以及轻量级内存队列。云托管队列最典型的代表是AWS SQS、阿里云消息队列RocketMQ、腾讯云CMQ。它们的优势是几乎零运维自动扩展提供死信队列和消息延迟等高级功能。如果你用的是某一家云厂商的函数计算直接选同一家的消息队列是最省心的——函数触发器可以无缝对接队列消息一来自动触发函数执行不用自己写拉取逻辑。Kafka和Pulsar这类自建消息系统适合已经有成熟集群、对消息吞吐量要求极高、需要消息有序性保证来支撑复杂事件流的场景。但在Serverless架构里自建Kafka有个尴尬点你为了跑一个消息队列还得维护一组服务器这就抵消了Serverless“无服务器”的优势。除非业务量已经大到云托管队列太贵否则我一般不建议前期就上Kafka。还有一种特殊情况只是函数内部自己重试几次或者一个请求拆成两三个子任务量很小值不值得引入外部队列这时候可以用函数计算的异步调用或者简单的内存队列比如Python的queue模块配合多线程。但注意内存队列只存在于单个执行环境里函数执行完就销毁了一旦进程被杀消息就丢。所以它只能做“进程内缓冲”不能做持久化消息。我个人的选型经验是默认选择云厂商的消息队列产品而不是自己搭。原因有三第一函数触发器集成度高配置一个事件源映射就完事第二托管队列天然具备高可用和持久化消息不会因为函数崩溃而丢失第三费用通常是按量计费微小型任务每月几乎可以忽略。2.3 事件源API网关、对象存储、数据库变更流异步处理不仅要处理“用户请求→后台任务”还要处理各种事件源的触发。Serverless的异步模式一大优势就是可以监听各种事件来一个事件就生成一条消息丢进队列。最常见的几个事件源API网关用户发来HTTP请求网关收到后立刻同步调用一个“入口函数”入口函数快速校验并往队列里塞一条消息然后返回202 Accepted。真正耗时的业务逻辑由另一个“处理函数”消费队列异步执行。这样用户感知到的接口延迟只有几十毫秒。对象存储比如用户上传了一个视频文件对象存储触发一个事件函数收到事件后把视频转码任务写入队列。转码可能需要几分钟完全在后台进行。数据库变更流比如数据库里插入了一条订单记录触发函数把订单数据同步到搜索服务或数据分析系统。这类操作本身就是异步的谁也不会让下单接口同步等待数据同步完成。定时触发器这就是本篇文章标题里提到的“serverless定时任务”场景。云厂商提供的定时触发功能可以按cron表达式每天、每小时触发一次函数。但这个“定时触发”本身是同步调用函数函数一次性执行。如果你的定时任务逻辑比较复杂建议在定时函数内部只做“拆分任务”的工作把每一条子任务丢到队列里再由消费者函数异步处理。我在做trae每日自动签到的时候就是这么设计的定时触发器每天早晨7点唤起一个函数这个“调度函数”读取需要签到的账号列表然后把每个账号的签到任务封装成一条消息写入队列。紧接着另一个“执行函数”被队列消息触发并发地处理每个账号的签到逻辑。这样即便账号数量增加我只需要调大队列消费者的并发数可以保证所有签到在几分钟内全部完成。3. 实战拆解一个定时任务的异步化改造3.1 场景每日自动签到为什么不能“硬等”先描述一下这个签到任务的业务背景trae是一个AI编程工具平台的样子遇到这种自动签到任务每天登录后可以领取一定奖励。我需要每天自动执行签到如果失败还要重试最好还能把签到结果发到通知渠道。这种定时任务看起来很简单一个cron触发器每天调用一个签到函数函数里遍历账号、调用签到接口、记录结果。问题在于签到不是只有一个账号。假设你有5个账号每个账号签到需要调用登录接口、获取token、再调用签到接口中间还可能遇到验证码或接口限流。正常情况一个账号耗时1秒但某个账号如果遇到网络抖动一次请求就可能卡20秒。5个账号串行执行最坏情况要等100秒直接超出函数超时上限。更麻烦的是签到接口偶尔会返回“系统繁忙”或“请求过于频繁”。如果你在函数里同步等待重试一旦函数超时被强制终止那这个账号的签到就漏了而且你没法确定它到底失败了没有。所以我把签到任务改成了异步模式定时器负责“拆单”队列负责“排队”消费函数负责“干每一单”。这样一来单个账号的签到失败只影响自己不会拖垮整批任务而且单个函数的执行时长也控制在了合理范围内。3.2 异步定时任务的标准架构整个架构分四层触发层、调度层、队列层、执行层。触发层云平台的定时触发器cron表达式每天07:00触发一次调度函数。调度层调度函数我们叫dispatcher读取配置列出所有待签到账号为每个账号生成一个“签到任务”消息发送到消息队列。调度函数本身执行非常快只做了读配置和发消息两件事毫秒级完成。队列层消息队列作为缓冲。每条消息包含账号ID、Cookie或令牌、签到类型、任务ID、时间戳等信息。执行层消息队列的触发器自动拉起“执行函数”函数一次处理一条消息完成签到并写入日志和状态存储。这种拆分的好处非常明显第一执行函数不再受“一次性处理所有账号”的时间限制每个函数只处理一个账号哪怕某个账号卡了30秒也只消耗一个函数实例的执行时间其他账号的消息还在队列里排队不会被耽误。第二失败重试变得很干净。如果某个账号签到失败执行函数抛异常消息队列会根据策略自动重试比如延迟1分钟再投递。第三扩展性好。账号增加到50个时队列消息量变多平台会自动拉起更多执行函数实例来消费消息不需要修改代码。我用的云平台消息队列和函数触发器的绑定方式是在函数配置里添加一个“消息队列触发器”指定队列名称和批量大小比如一次最多拉取10条消息。这样队列里一旦有消息平台就会自动调用执行函数并把消息体作为事件参数传入。3.3 伪代码与关键配置从触发到执行下面是我当时写的简化伪代码逻辑非常直白重点是看异步模式怎么流转。# dispatcher: 每日定时触发的调度函数 import json import os from qcloud import send_message # 假设的云厂商SDK def main_handler(event, context): accounts get_signed_accounts() # 从配置存储读取账号列表 task_id generate_task_id(daily_sign, date_today()) for account in accounts: message { task_id: task_id, account_id: account[id], cookie: account[cookie], platform: account[platform], retry_count: 0 } send_message(queue_namesign_task_queue, messagejson.dumps(message)) # 记录一次调度日志方便追踪 log_metric(dispatched_tasks, len(accounts)) return {dispatched: len(accounts)}# worker: 消费队列中的签到任务 import json import time def main_handler(event, context): # 云平台的触发器会把消息列表传入 event records event.get(Records, []) for record in records: msg_body json.loads(record[msgBody]) account_id msg_body[account_id] cookie msg_body[cookie] try: result do_sign_in(account_id, cookie) save_sign_result(msg_body[task_id], account_id, result) except SignInRateLimited as e: # 主动抛出异常让平台重试 raise e except Exception as e: log_error(sign_failed, account_id, e) # 也可以捕获后标记失败结果写入状态表避免重复重试 save_sign_result(msg_body[task_id], account_id, {status: failed, reason: str(e)})配置上需要注意几个点定时表达式云厂商的cron表达式一般会标注时区一定要看清楚。我此前吃过亏默认用的是UTC时区导致定时任务比北京时间晚了8小时。配置cron时要么显式指定时区要么在函数里做时区换算。消息队列的可见性超时如果执行函数正在处理一条消息但还没处理完队列的“可见性超时”就到了这条消息会被再次投递给其他函数实例。可见性超时应该设置得比你函数最大执行时间稍长一些。比如函数最大执行时间是60秒可见性超时就设成90秒。队列的消息保留时间云厂商默认一般保留1-4天。如果下游故障太久消息可能过期。注意提前设置合理的保留时间并配置死信队列避免消息在故障期间丢了。这整套流程从定时触发到消息入队再到函数消费真正做到了“入口轻、执行稳、可重试”。4. 幂等性与重试异步系统最容易翻车的两个点4.1 幂等设计重复执行不等于重复扣款异步系统里消息至少会被投递一次这是分布式系统的硬性现实。即使你用的是云托管队列也只能保证“不丢消息”很难保证“完全不重复”。尤其是上面提到的“可见性超时”机制如果函数执行超过超时时间消息会被投递给另一个实例这时同一个任务就有两个实例在同时处理。如果你在同步服务里写业务逻辑天然不需要考虑同一个请求被处理两次的情况因为客户端发送一次你处理一次。但异步系统不同你必须假设每一条消息都可能被处理多次而且可能是并行的。幂等性设计通俗来说就是无论这个任务被重复执行多少次最终产生的状态应该和执行一次完全一样。还是拿签到举例。假设签到成功后要增加用户积分。如果签到接口本身不是幂等的你重复调用两次就会加两次积分。那么你的函数应该在调用第三方签到接口之前先检查一下本地记录这个task_id和account_id组合是不是已经处理过了如果处理过直接返回“重复消息已忽略”。更通用的做法是引入“去重表”或者“状态表”。用task_id作为唯一键在处理前插入一条statuspending的记录处理完成后把status更新为success或failed。如果发现task_id已经存在且是success直接丢掉这条消息。这相当于给异步任务加了一堵“回执墙”。当然有些场景天然幂等比如“设置用户昵称为xxx”不管执行几次结果都一样。但涉及“增加余额”“发送短信”“创建订单”这类操作几乎是必须做幂等控制的。4.2 重试策略指数退避和死信队列异步任务失败后最自然的处理方式就是重试。但重试不是无脑循环如果下游服务已经处于过载状态你疯狂重试只会加重故障。标准做法是指数退避第一次失败后等待1秒重试第二次等待2秒第三次4秒最多等待30秒或更久同时每次重试之间增加一些随机抖动避免多个消息同时重试导致“重试风暴”。云厂商的消息队列大多支持配置“消息重试策略”比如最大重试次数3次每次间隔递增。我在trae签到项目里把最大重试次数设成了5次间隔从10秒到5分钟逐步扩大。因为这个签到任务每天只有一次机会如果当天早上没签上当天就错过了所以多试几次值得。但如果是非关键的日志处理任务重试3次就够了。重试次数耗尽后消息应该进入死信队列。死信队列的作用是“收容所”专门存放反复处理失败的消息。你可以在死信队列里设置另外一个消费者函数专门用来告警或者人工介入处理。我一般会给死信队列挂一个告警一旦有消息落入死信立刻发短信或钉钉通知这样我能在当天发现签到任务挂掉了而不是月底才发现漏签了半个月。4.3 至少一次 vs 至多一次消息队列提供不同的投递语义你要根据业务容忍度来选。至多一次At most once消息可能被丢失但不会重复。适合容忍丢数据、不能容忍重复的场景比如实时日志采集偶尔丢几条日志无所谓。至少一次At least once消息不丢但可能重复。这是大多数云厂商消息队列的默认行为。适合签到、订单处理这类“宁可重复也不能漏”的场景。精确一次Exactly once看起来完美但实现代价高通常需要业务幂等配合或依赖下游唯一约束。现实里很少有一个分布式消息系统能完全不加业务配合就做到精确一次。我的经验是不要追求绝对精确一次而是“至少一次业务幂等”组合。消息队列保证消息不会丢最多重投业务层用幂等去重兜底这样既简单又可靠。这就像快递物流公司保证包裹不弄丢但是可能送两次你家的门卫大爷用登记本记录包裹编号第二次送来的重复包裹直接拒收。5. 冷启动、并发和成本异步模式下的三座山5.1 冷启动对异步任务的影响冷启动是老生常谈的话题。异步处理模式下冷启动带来的“延迟”比较特殊你从用户请求的角度看入口函数总是热着的因为它每天频繁被调用冷启动几乎可忽略。而从队列消费函数的角度看如果消息不是持续到达队列消费者可能长期空闲那么每当一批新消息到来消费函数都可能是冷启动状态要等待运行时初始化几百毫秒到几秒不等。这个延迟在大部分异步场景下是可接受的。比如定时签到任务某一批消息到达后执行函数花2秒冷启动也无所谓因为用户感知不到。但如果你是在处理实时性要求极高的交互式异步请求比如用户上传图片后希望稍后推送给客户端冷启动会导致推送延迟增加。这时可以通过预留并发实例或配置定时预热来缓解。我还踩过一个坑消费函数如果同时被大量消息触发平台会创建N个并发实例这N个实例全部冷启动瞬间把下游接口流量打满。所以要注意控制消费者的最大并发数。有的云平台支持配置函数的“单实例并发度”即每个实例同时处理几条消息。如果为了减少冷启动可以把单实例并发度调大比如一个实例同时消费5条消息这样实例数少了冷启动也少了。5.2 并发限制与削峰填谷每个云账号的函数并发配额都是有限的。比如某个账号默认并发配额是100意思是同时最多跑100个函数实例。如果你的异步队列里堆积了1万条消息消费函数瞬间拉起100个实例去处理虽然处理得快但可能触发你的下游接口限流或者把数据库连接打爆。所以异步模式下的并发不是越大越好而是要有节制。消息队列的消费者应该像调节水龙头一样根据下游承载能力来调节流量。具体手段一是设置函数的并发上限二是通过消息队列的“批量拉取大小”和“拉取间隔”来控制速率三是如果下游服务有明确的QPS限制可以在消费函数内部做一个简单的限流器通过令牌桶算法在本地控制调用速率。削峰填谷的原理就是消息队列缓冲了高峰期的全部请求消费者按照下游可承受的速率慢慢消费把原来“瞬时高峰”拉平成了“平稳长坡”。这样即使你的签到任务在早上7点瞬间产生大量消息数据库也不会被冲垮因为消费函数以固定速率在那里一条条地处理。5.3 成本核算异步真的省钱吗很多人误以为异步处理会增加成本因为多了一个消息队列和多个函数调用。实际上异步模式通常更省钱关键要算清楚账。先看函数计费。Serverless是按调用次数和资源使用时长收费的。如果同步处理入口函数和执行函数都在一次调用里时长是两者的总和。异步化之后入口函数只做了快速校验和发消息时长从几秒降到了几十毫秒执行函数还是消耗相同的时长。总费用来说入口函数变便宜了执行函数没变所以整体费用下降。再看消息队列计费。大多数云厂商的消息队列有免费额度例如每月前100万条消息免费超出后按百万条计费。对于绝大多数中小业务来说这部分费用几乎可以忽略不计。还有隐性的成本如果同步函数经常超时被强制终止已经做完的数据库操作可能会产生脏数据清理脏数据的人工成本往往比云服务费用贵得多。异步方式让你可以精确控制重试和状态少踩坑就是省钱。从我实测的数据来看同样的签到任务同步一次性处理50个账号函数时长累计可能超过70秒计费时间很长异步拆分成50条消息每条消息平均处理1秒调度函数只耗时0.05秒执行函数的总计费时长也是50秒。看似没差多少但同步模式下某一次失败可能就要人工处理而异步模式自动重试成功省下的运维精力是无法量化的。6. 常见问题与排查技巧实录6.1 任务“消失”了异步系统里用户最常反馈的问题就是“我明明提交了怎么没后续”。排查思路如下第一检查入口函数是否成功往队列里发了消息。看日志里有没有dispatch成功记录消息队列的控制台里有没有消息堆积指标。如果入口函数执行成功但是没有消息那基本是SDK调用问题检查队列名称、区域、权限。第二检查消费函数是否被触发。看消费函数的调用次数是不是0。如果队列有消息但消费函数没被触发多半是函数触发器没有正确关联队列或者队列事件的权限配置有问题。翻一下云平台的触发器配置确认“消息队列触发器”是启用状态。第三检查死信队列。如果消费函数处理消息失败而且重试耗尽消息会进入死信队列。去死信队列里翻一翻往往能在那里找到“消失”的任务。这基本能覆盖99%的“消息消失”现象。6.2 消息重复导致的脏数据如果业务里出现了重复记录先别急着怀疑消息队列“数据错乱”大概率是可见性超时设短了。举一个实例消费函数里调用了第三方签到接口第三方接口响应很慢函数执行了80秒但队列的可见性超时只有60秒。60秒后消息被队列重新投递给另一个实例。此时第一个实例可能还在等签到接口响应第二个实例又开始调同一个签到接口这就会重复签到可能造成积分重复发放。解决办法有三一是把可见性超时调大到函数执行上限的1.5倍以上二是在代码里加分布式锁或去重表三是让签到接口本身配套幂等键。三个措施里最省事的是调整可见性超时但也是最容易忽视的。我建议你在设计新函数时就把“可见性超时 函数超时时间 30秒”作为默认规则写进检查清单。6.3 定时任务偶发超时定时任务触发只是“定时器到点调用函数”但调度函数本身也可能偶发超时原因多半是调度逻辑写了耗时的操作。我早期犯过这样的错误在调度函数里顺便做了数据统计统计完再发消息结果数据处理慢整个调度函数超时了。后来我把调度函数改成“只读取账号列表、只发消息”所有重活全部丢给执行函数。现在不管账号数量怎么涨调度函数永远在1秒内完成。所以定时任务异步化的一个细节法则是定时函数里的代码越少越好宁可多一个步骤也不要拖慢触发链路。你可以把“准备数据”和“触发任务”拆成两个动作先由定时函数把数据准备好放在存储里再由另一个函数去读存储、发消息。没必要挤在一个函数里。6.4 日志与可观测性异步系统的链路跨越了定时器、调度函数、队列、执行函数多个环节一旦出错最怕的就是查不到日志。我的建议是建立一套“单请求全链路ID”的日志体系。在每个定时任务生成时创建一个task_id然后在所有函数日志里带上这个ID。无论是调度函数日志、队列消息内容、执行函数日志都打印task_id。这样排查问题时只需要在日志系统里按task_id搜索就能找到同一批任务的全部执行过程。另外要给关键环节埋点任务调度数、入队消息数、消费成功数、消费失败数、死信数。这些指标可以直接在云监控里查看。我建议设置两个必看的告警一个是“死信队列有消息”另一个是“消费函数失败率超过阈值”。有了这两个告警你基本能在用户发现前自己先发现问题。我自己的trae每日自动签到项目就是靠这套日志体系救回来的。有一次trae那边改了接口返回结构签到脚本大面积失败消息进入死信队列监控立刻发了告警。我当天早上8点就看到了消息在用户还没反应过来之前已经改好代码把死信队列里的任务重新放回主队列消费当天签到没有漏掉。这要是没有异步模式和配套监控很可能要到月底统计奖励时才发现问题那时候已经没法补救了。异步处理模式说到底是Serverless的精髓所在。它逼着我们用“事件驱动”的思维来设计系统把长任务拆小把失败隔离把削峰做掉。虽然前期要多写一些幂等和重试的代码但换来的是系统的稳定性和省心程度。如果你也正打算做定时任务批量处理类的项目我强烈建议一开始就按异步的骨架来搭别等超时和丢消息了再回头改。那个时候的痛苦绝不是多写两行代码能弥补的。
返回列表