面向Agent系统的Java后端知识总览(中)

Agent 后端进入稳定运行阶段后,核心问题会从“能否执行任务”转向“任务能否可靠完成”。一次任务可能持续数分钟,中间经历模型调用、知识检索、业务工具执行和结果校验。服务重启、消息重复投递、缓存失效或数据库写入失败,都可能导致任务状态混乱、重复执行和结果丢失。

MySQL、Redis 和消息队列分别承担可靠数据、临时状态与异步任务的管理。三者之间需要明确边界。MySQL 保存最终可信的业务事实,Redis 提供高频访问和流量控制,消息队列负责传递任务与故障重试。

目录

一、任务状态的持久化设计

二、Redis的状态管理边界

三、Redis限流与分布式锁

四、消息队列的长任务调度

五、Kafka、RabbitMQ与RocketMQ的选择

六、数据库与消息的一致性

七、重试、超时与失败分类

八、数据中间件的运行风险


一、任务状态的持久化设计

Agent 任务通常具有较长的生命周期,状态变化应当保存在关系数据库中。服务实例重启或扩容后,任务记录仍然可以恢复和查询。一个基础任务表至少需要包含任务编号、用户编号、执行状态、请求摘要、结果位置、失败原因和时间字段。

CREATE TABLE agent_task ( id VARCHAR(64) PRIMARY KEY, user_id BIGINT NOT NULL, status VARCHAR(32) NOT NULL, request_hash VARCHAR(64) NOT NULL, result_url VARCHAR(500), error_message VARCHAR(500), retry_count INT NOT NULL DEFAULT 0, created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL, UNIQUE KEY uk_user_request (user_id, request_hash), KEY idx_status_updated (status, updated_at) );

uk_user_request可以阻止同一用户在业务允许的范围内重复创建相同任务,idx_status_updated适合扫描长时间停留在执行状态的任务。索引应围绕实际查询条件设计,避免为每个字段单独建立索引。索引数量增加会提高查询能力,也会带来额外的写入和存储成本。

任务状态更新应当具备方向约束。例如,只有处于ACCEPTED状态的任务才能进入RUNNING,已经成功的任务不能被迟到的失败消息覆盖

UPDATE agent_task SET status = 'RUNNING', updated_at = NOW() WHERE id = ? AND status = 'ACCEPTED';

Java 代码通过更新行数判断状态是否成功迁移。

int updated = taskRepository.markRunning(taskId); if (updated == 0) { return; }

这种条件更新同时具备状态校验和并发控制能力。多个消费者同时处理相同任务时,只有一个消费者能够成功修改状态。其余消费者读取到更新行数为零后结束处理,不需要额外增加分布式锁。

MySQL InnoDB支持多版本并发控制和行级锁,并提供四种标准事务隔离级别事务范围应限制在必要的数据库操作中,模型调用、文件处理和远程接口请求不应包含在长事务内,否则会持续占用数据库连接并扩大锁冲突范围。(MySQL开发者区)


二、Redis的状态管理边界

Redis适合保存变化频繁、允许过期和读取量较高的数据,例如会话上下文、任务进度、临时授权信息、模型配额和热点配置。数据库中的任务状态仍然是最终可信依据,Redis中的数据可以在失效后重新构建。

任务执行过程中可以将进度写入 Redis,并设置合理的过期时间。

String key = "task:progress:" + taskId; redisTemplate.opsForHash().put(key, "step", "TOOL_CALL"); redisTemplate.opsForHash().put(key, "percent", "60"); redisTemplate.expire(key, Duration.ofHours(2));

进度数据适合高频更新,没有必要让每一次百分比变化都写入 MySQL。任务完成后,再将最终状态和结果位置持久化到数据库。客户端查询时可以优先读取 Redis,缓存不存在时回退到 MySQL。

缓存设计需要防止穿透、击穿和雪崩。不存在的数据可以短时间缓存空结果,热点数据可以通过互斥加载或逻辑过期避免大量请求同时访问数据库,不同缓存键应设置一定的过期时间抖动,降低大量键同时失效的风险。


三、Redis限流与分布式锁

模型调用通常受制于接口配额、并发连接和推理资源。单个用户或租户如果在短时间内提交大量任务,可能占满执行线程和模型配额。Redis提供原子计数与键过期能力,可以实现固定窗口、滑动窗口和令牌桶等限流方式。(Redis)

固定窗口限流的核心逻辑较为简单。

public boolean allow(String userId) { String key = "rate:task:" + userId; Long count = redisTemplate.opsForValue().increment(key); if (count != null && count == 1) { redisTemplate.expire(key, Duration.ofMinutes(1)); } return count != null && count <= 20; }

这段代码限制单个用户每分钟最多创建 20 个任务,适合基础保护。固定窗口在时间边界附近可能允许短时间突发请求,要求更严格时可以使用 Lua 脚本实现滑动窗口或令牌桶,将多个 Redis 操作合并为一次原子执行

分布式锁适合协调跨实例的稀缺资源,例如防止多个节点同时刷新同一份配置。获取锁时需要同时设置唯一值和过期时间,释放锁时必须验证锁的所有者。

Boolean locked = redisTemplate.opsForValue().setIfAbsent( lockKey, requestId, Duration.ofSeconds(30) );

锁的有效期过短可能在业务尚未完成时提前释放,过长则会延迟故障恢复。主从切换、网络暂停和锁续期也会影响锁的安全性。Redis官方将分布式锁作为需要专门算法处理的协调模式,简单的SETNX只能覆盖有限场景。(Redis)

对于任务幂等和状态流转,数据库唯一约束与条件更新通常更加直接。分布式锁应当用于确实需要跨节点互斥的操作,避免成为所有并发问题的默认答案。


四、消息队列的长任务调度

消息队列将任务接收和任务执行分离。接口创建数据库记录后发送消息(生产者),消费者读取消息并执行模型、检索和工具调用。任务执行速度发生波动时,消息可以在队列中等待

public TaskResponse create(TaskRequest request) { AgentTask task = taskRepository.save( AgentTask.accepted(request) ); messageProducer.send( new TaskMessage(task.getId()) ); return new TaskResponse( task.getId(), task.getStatus() ); }

消费者不能假设一条消息只会到达一次。网络超时、确认丢失和消费者重启都可能导致重新投递,因此消费逻辑必须具备幂等性。

public void consume(TaskMessage message) { int updated = taskRepository.markRunning( message.taskId() ); if (updated == 0) { return; } try { TaskResult result = taskExecutor.execute( message.taskId() ); taskRepository.markSucceeded( message.taskId(), result.location() ); } catch (Exception exception) { taskRepository.markFailed( message.taskId(), exception.getMessage() ); throw exception; } }

消息确认应当发生在业务处理成功之后。业务执行失败时,消息可以进入重试流程。持续失败的消息需要进入死信队列,避免单条异常任务无限占用消费者资源。


五、Kafka、RabbitMQ与RocketMQ的选择

三种消息中间件都可以传递异步任务,设计侧重点存在差异。

中间件主要特点常见应用
Kafka分区日志、高吞吐、消息保留事件流、审计日志、数据管道
RabbitMQ路由灵活、确认机制清晰工作队列、业务消息、复杂路由
RocketMQ重试、延迟消息和顺序消息支持完善长任务、延迟任务、业务事件

Kafka生产者可以启用幂等写入,降低生产者重试造成的重复记录。启用幂等性时,需要配合acks=all、重试配置和连接内请求数量约束。Kafka的生产者幂等只能解决消息写入 Kafka 时的部分重复问题,消费者执行数据库操作和外部调用时仍然需要业务幂等。(Kafka)

RabbitMQ使用 Publisher Confirm确认消息是否被 Broker 接收,使用 Consumer Acknowledgement确认消费者是否完成处理。两种确认分别覆盖生产和消费方向,无法替代业务层面的状态判断。(RabbitMQ)

RocketMQ能够根据消费重试策略重新投递失败消息,超过最大重试次数后将消息发送到死信队列。死信消息需要配套查询、告警和人工恢复机制,否则仍然会形成不可见的任务积压。(RocketMQ)

中间件选择需要结合已有技术栈、运维能力、消息规模和业务语义。普通任务平台通常不需要同时引入多种消息队列,统一消息模型和故障处理方式更有利于长期维护。


六、数据库与消息的一致性

创建任务时通常需要完成两项操作,一项是向 MySQL 写入任务记录,另一项是向消息队列发送任务消息。如果数据库提交成功而消息发送失败,任务会永久停留在已接收状态;如果消息发送成功而数据库回滚,消费者又会找不到对应任务。

本地消息表可以降低这种双写风险。业务事务中同时写入任务表和消息表,再由独立程序扫描待发送消息。

BEGIN; INSERT INTO agent_task ( id, user_id, status, request_hash, created_at, updated_at ) VALUES (?, ?, 'ACCEPTED', ?, NOW(), NOW()); INSERT INTO task_outbox ( event_id, task_id, status, created_at ) VALUES (?, ?, 'PENDING', NOW()); COMMIT;

后台发布器读取PENDING消息,发送成功后将其更新为SENT。消息表可能被重复扫描,因此生产者和消费者仍然需要幂等控制。这种方式通过数据库本地事务保证任务记录与待发送事件同时存在,避免依赖跨中间件事务。

对于规模较小的系统,也可以使用定时任务扫描长时间停留在ACCEPTED状态的任务并重新投递。该方案实现简单,但补偿存在时间延迟,适合作为恢复机制,不宜完全替代正常消息投递流程。


七、重试、超时与失败分类

所有失败都进行重试会放大系统压力网络抖动、临时限流和服务短暂不可用通常可以重试参数错误、权限不足和业务规则冲突则应直接失败。模型接口返回超时或限流时,可以采用指数退避,并限制最大重试次数。

long delaySeconds = Math.min( 60, 1L << retryCount );

重试前需要判断当前任务状态,已经成功、取消或超过截止时间的任务不应再次执行。涉及外部扣费、通知发送和业务写入时,还应携带稳定的幂等键,使下游能够识别重复请求。

失败记录至少应保留错误分类、重试次数、最近执行时间和原始任务编号。错误信息只保存必要内容,避免将访问令牌、用户隐私和完整模型输入直接写入日志或数据库。


八、数据中间件的运行风险

MySQL需要重点观察慢查询、连接池耗尽、锁等待和大事务。Redis需要关注内存使用、热键、过期键集中删除和缓存命中率。消息队列需要关注生产失败、消费积压、重复消费、重试次数和死信数量。

任务执行系统不能只监控接口是否存活。数据库中的各状态任务数量、Redis限流触发次数、消息积压长度、任务平均等待时间和死信数量,更能反映系统是否处于健康状态。RocketMQ官方指标中也包含消费结果和死信消息数量,说明重试与死信本身就是消息系统需要持续观测的运行状态。(RocketMQ)

MySQL保存业务事实,Redis管理短期状态和访问速度,消息队列承接异步任务和失败恢复。三者通过任务编号、状态机和幂等约束连接后,Agent后端才能在服务重启、消息重投和流量波动中保持稳定。