Dubbo框架服务间RPC接口调用超时处理方案
问题背景:
在我的微服务中,有A服务做定时任务,B服务做数据收集。服务间通过dubbo框架进行通信。 现在我用A服务通过dubbo的rpc接口调用B服务的数据收集任务。因为任务耗时久的原因,A服务会报错:
Timeout: 30000ms
在此基础上,有两个思路进行优化
- 在调用接口时添加超时设置
▼java复制代码@DubboReference( timeout = 7200000, // 2小时超时,根据实际需要调整 )
- 将同步调用改为异步调用
因为我的服务端方法是一个数据收集解析的长时间任务,更具数据量大小,最大执行时间可能达到小时级别。所以第一种方案被我弃掉。采用了异步调用的改造方案。
一.异步改造
服务端改造:
在原有代码基础上添加@Async注解
▼java复制代码@Override @Async // 使用Spring的异步执行 public CompletableFuture<String> addDatas() { try { log.info("开始执行长时间数据收集任务..."); dataService.addDatas(); log.info("数据收集任务完成"); } catch (Exception e) { log.error("数据收集任务失败", e); } }
客户端改造:
客户端则需要在@DubboReference注解中指定 async = true开启异步调用
▼java复制代码@DubboReference( async = true, timeout = 7200000, // 2小时超时,根据实际需要调整 ) private DataRpcService dataRpcService; @Override public void run() { log.info("定时任务 - 数据收集"); // 发起异步调用 dataRpcService.addDatas(); }
这个时候,我们的服务就不会出现接口超时的报错了。但是这样会引发一个日志记录错误的问题。就是任务日志会在我们发起异步调用后,直接记录结束。即使这时服务端还在执行任务。为了解决这个问题。需要引入一个核心组件CompletableFuture。
一.核心组件CompletableFuture,解决日志记录不准问题
通过与AI交互。发现了Dubbo框架支持基于CompletableFuture接口的异步调用。而不会阻塞当前线程。该Future对象相当于一个"凭证",消费者可以:
工作原理:
- 立即返回:A服务发起调用后,Dubbo框架会立即构造并返回一个
CompletableFuture对象,而不会阻塞等待B服务的业务逻辑执行完毕。这时,A服务发起的这次RPC(远程过程调用)在网络层面可以很快结束,避免了因长时间等待而触发的超时。 - 后台执行:B服务在接收到请求后,开始执行实际耗时很长的
addDatas方法。这个执行过程与A服务已经解耦。 - 结果回调:当B服务的方法执行完毕后,Dubbo框架会隐式地将结果设置回之前A服务拿到的那个
CompletableFuture对象。此时,在A服务中,通过future.get()等待这个结果,或者通过future.whenComplete()添加的回调函数就会被触发,从而获取到最终的执行结果("SUCCESS"或错误信息)。
服务端接口实现:
▼java复制代码@Override @Async // 使用Spring的异步执行 public CompletableFuture<String> addDatas() { try { log.info("开始执行长时间数据收集任务..."); dataService.addDatas(); log.info("数据收集任务完成"); return CompletableFuture.completedFuture("SUCCESS"); } catch (Exception e) { log.error("数据收集任务失败", e); return CompletableFuture.completedFuture("FAILED: " + e.getMessage()); } }
客户端实现:
CompletableFuture提供了future.get()和future.whenComplete()两种方法去处理。
同步等待机制:future.get()会阻塞当前线程,直到B服务中的addData方法执行完成并返回结果。这使得能够准确记录从任务开始到结束的完整时间段。
▼java复制代码@DubboReference( async = true, timeout = 7200000, // 2小时超时,根据实际需要调整 retries = 0 // 异步调用不重试 ) private PdcUserLogRpcService pdcUserLogRpcService; @Override public void run() { log.info("定时任务 - UserLog 数据收集"); // 发起异步调用 pdcUserLogRpcService.addPdcUserLog(); CompletableFuture<String> future = RpcContext.getContext().getCompletableFuture(); try { // 同步等待异步任务完成(会阻塞当前线程,但不会超时) String result = future.get(7100, TimeUnit.SECONDS); // 略小于2小时,给日志记录留时间 if ("SUCCESS".equals(result)) { log.info("数据收集任务执行成功"); } else { log.warn("数据收集任务完成但有异常: {}", result); // 抛出异常让LoggableJobWrapper捕获并记录失败状态 throw new RuntimeException("远程任务返回失败: " + result); } } catch (Exception e) { log.error("任务失败", e); throw new RuntimeException("任务失败", e); } }
回调函数机制:future.whenComplete()会在任务结束后被唤醒,并拿到返回值。便于后续操作。
▼java复制代码@DubboReference(async = true, timeout = 7200000, retries = 0) private DataRpcService dataRpcService; // 用于跟踪长时间运行的任务 private final Map<String, CompletableFuture<String>> runningTasks = new ConcurrentHashMap<>(); @Override public void run() { String taskId = "DataTask-" + System.currentTimeMillis(); log.info("开始异步调用数据收集任务, 任务ID: {}", taskId); dataRpcService.addDatas(); CompletableFuture<String> future = RpcContext.getContext().getCompletableFuture(); // 保存任务引用 runningTasks.put(taskId, future); future.whenComplete((result, exception) -> { // 任务完成,从运行列表中移除 runningTasks.remove(taskId); if (exception != null) { log.error("任务 {} 执行失败", taskId, exception); } else { log.info("任务 {} 完成,结果: {}", taskId, result); } }); log.info("任务 {} 已提交,当前运行中的任务数量: {}", taskId, runningTasks.size()); } // 获取当前运行任务状态的方法 public int getRunningTaskCount() { return runningTasks.size(); }
评论
问答助学
相关内容
0个评论
全部评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
