当前位置: 代码网 > it编程>编程语言>Java > Java实现多步骤异步任务编排的三种主流方案

Java实现多步骤异步任务编排的三种主流方案

2026年08月26日 Java 我要评论
背景在仓储物流、产线自动化等场景中,我们经常需要通过 java 后端协调 agv 小车完成一系列连续动作:先通过 http 接口下发移动任务,等待小车到达指定位置后,再通过 opc ua 协议控制机械

背景

在仓储物流、产线自动化等场景中,我们经常需要通过 java 后端协调 agv 小车完成一系列连续动作:先通过 http 接口下发移动任务,等待小车到达指定位置后,再通过 opc ua 协议控制机械臂完成夹取,最后再下发移动任务到目的地。

这类场景的核心难点在于:http 任务下发是异步的,opc ua 指令执行也是异步的,而多个步骤之间存在严格的先后依赖关系。本文将介绍三种主流的实现方案,从轻量级到重量级依次展开。

问题拆解

一个典型的 agv 搬运流程如下:

  1. 通过 http 接口向 rcs(机器人调度系统)下发移动任务,目标为取货点
  2. 监听任务状态,等待小车到达取货点
  3. 到达后,通过 opc ua 协议向机械臂下发夹取指令
  4. 监听夹取指令执行结果,等待夹取完成
  5. 夹取完成后,再次通过 http 接口下发移动任务,目标为目的地
  6. 监听任务状态,等待小车到达目的地,流程结束

可以看到,整个流程是一个典型的多步骤异步流程编排问题,每一步都需要等待上一步完成后再触发。

方案一:completablefuture 链式编排(推荐,轻量级)

java 8 引入的 completablefuture 天然适合这种"步骤 a 完成 → 触发步骤 b → 触发步骤 c"的场景。通过 thencompose 方法可以优雅地串联有依赖关系的异步操作,避免嵌套回调(即"回调地狱")。

@service
public class agvtaskorchestrator {

    private static final logger log = loggerfactory.getlogger(agvtaskorchestrator.class);

    @autowired
    private agvhttpservice agvhttpservice;
    @autowired
    private opcuaservice opcuaservice;
    @autowired
    private taskstatuslistener statuslistener;

    public completablefuture<void> executepickandmove(string agvid, string targetstation) {

        // 步骤1:下发http任务,前往取货点,等待到达
        completablefuture<string> gotopoint1 = agvhttpservice.sendmovetask(agvid, "pick_point")
                .thencompose(taskid -> statuslistener.waitforcomplete(taskid));

        // 步骤2:到达后,通过opc ua下发夹取指令,等待夹取完成
        completablefuture<void> grabaction = gotopoint1
                .thencompose(arrivedtaskid -> opcuaservice.sendgrippercommand(agvid, "grab"))
                .thencompose(cmdid -> opcuaservice.waitforcommandcomplete(cmdid));

        // 步骤3:夹取完成后,下发移动任务到目的地,等待到达
        completablefuture<string> gototarget = grabaction
                .thencompose(v -> agvhttpservice.sendmovetask(agvid, targetstation))
                .thencompose(taskid -> statuslistener.waitforcomplete(taskid));

        return gototarget.thenaccept(v -> log.info("agv[{}] 全流程执行完成", agvid));
    }
}

关键点thencompose 用于串联有依赖关系的异步操作,它会将前一步的返回结果传递给下一步,同时避免 completablefuture<completablefuture<...>> 的嵌套问题。

适用场景:步骤较少(3~5 步)、逻辑简单、没有复杂分支的场景。

方案二:状态机模式(适合复杂流程)

当流程步骤增多、出现分支逻辑(如夹取失败要重试、异常要回退)时,completablefuture 链式编排会变得难以维护。此时推荐使用状态机模式来管理流程状态。

public enum agvflowstate {
    idle,                // 空闲
    moving_to_pick,      // 前往取货点
    grabbing,            // 夹取中
    moving_to_target,    // 前往目的地
    completed,           // 完成
    failed               // 失败
}

@service
public class agvstatemachine {

    private static final logger log = loggerfactory.getlogger(agvstatemachine.class);

    private final map<string, agvflowstate> taskstates = new concurrenthashmap<>();

    @autowired
    private agvhttpservice agvhttpservice;
    @autowired
    private opcuaservice opcuaservice;

    /**
     * 收到http任务完成回调时调用
     */
    public void ontaskarrived(string taskid) {
        agvflowstate currentstate = taskstates.get(taskid);
        if (currentstate == null) return;

        switch (currentstate) {
            case moving_to_pick:
                log.info("agv已到达取货点,下发夹取指令, taskid={}", taskid);
                taskstates.put(taskid, agvflowstate.grabbing);
                opcuaservice.sendgrippercommand(taskid, "grab");
                break;
            case moving_to_target:
                log.info("agv已到达目的地,流程结束, taskid={}", taskid);
                taskstates.put(taskid, agvflowstate.completed);
                break;
            default:
                log.warn("收到意外的任务完成回调, taskid={}, state={}", taskid, currentstate);
        }
    }

    /**
     * 收到opc ua指令完成回调时调用
     */
    public void ongrippercomplete(string taskid) {
        agvflowstate currentstate = taskstates.get(taskid);
        if (currentstate == agvflowstate.grabbing) {
            log.info("夹取完成,下发移动任务到目的地, taskid={}", taskid);
            taskstates.put(taskid, agvflowstate.moving_to_target);
            agvhttpservice.sendmovetask(taskid, "target_station");
        }
    }

    /**
     * 启动一个新的搬运流程
     */
    public void startflow(string taskid) {
        taskstates.put(taskid, agvflowstate.moving_to_pick);
        agvhttpservice.sendmovetask(taskid, "pick_point");
    }
}

优势:每个步骤的状态转换清晰可控,方便加日志、异常处理和断点恢复。当流程变复杂时,还可以引入 spring statemachine 框架来进一步简化状态定义和转换规则。

适用场景:步骤较多、有分支/重试/回退逻辑的复杂场景。

方案三:回调 + 事件驱动(webhook 模式)

如果你的 rcs 系统支持 webhook 回调(如海康 rcs-2000 等主流系统),可以用事件驱动的方式来实现,让 rcs 在任务状态变更时主动推送通知。

@restcontroller
@requestmapping("/agv")
public class agvcallbackcontroller {

    private static final logger log = loggerfactory.getlogger(agvcallbackcontroller.class);

    @autowired
    private agvtaskorchestrator orchestrator;

    /**
     * rcs回调接口:任务状态变更时主动推送
     */
    @postmapping("/callback")
    public responseentity<?> onagvcallback(@requestbody agvcallbackdto callback) {
        log.info("收到agv回调, taskid={}, status={}", callback.gettaskid(), callback.getstatus());

        switch (callback.getstatus()) {
            case "arrived":
                orchestrator.onarrived(callback.gettaskid());
                break;
            case "completed":
                orchestrator.onmovecomplete(callback.gettaskid());
                break;
            case "failed":
                orchestrator.onfailed(callback.gettaskid(), callback.geterrormessage());
                break;
            default:
                log.warn("未知的回调状态: {}", callback.getstatus());
        }
        return responseentity.ok().build();
    }
}

适用场景:rcs 系统支持回调/webhook 推送的场景,实时性最好,无需轮询。

关于"监听任务完成"的两种实现方式

在上述方案中,"监听任务完成"是一个关键环节,通常有两种实现方式:

方式实现思路适用场景
轮询定时调用查询接口检查任务状态rcs 不支持回调时
回调/webhook暴露 http 接口,rcs 主动推送状态变更rcs 支持回调时(推荐)

轮询方式的实现示例:

public completablefuture<string> waitforcomplete(string taskid) {
    return completablefuture.supplyasync(() -> {
        int maxretries = 300; // 最多等待5分钟
        int retrycount = 0;

        while (retrycount < maxretries) {
            taskstatus status = agvhttpservice.querytaskstatus(taskid);
            if (status == taskstatus.completed) {
                return taskid;
            }
            if (status == taskstatus.failed) {
                throw new runtimeexception("任务执行失败, taskid=" + taskid);
            }
            try {
                thread.sleep(1000); // 每秒轮询一次
            } catch (interruptedexception e) {
                thread.currentthread().interrupt();
                throw new runtimeexception("等待任务完成被中断", e);
            }
            retrycount++;
        }
        throw new runtimeexception("任务超时, taskid=" + taskid);
    });
}

注意:轮询方式一定要设置超时机制,避免无限等待导致线程泄漏。

方案对比与选型建议

维度completablefuture状态机事件驱动
复杂度
可维护性步骤多时较差
实时性取决于监听方式取决于监听方式最好
异常处理需手动处理天然支持需手动处理
适用场景3~5步简单流程多步骤复杂流程实时性要求高

实际选型建议

  • 步骤少(3~5 步)、逻辑简单:用 completablefuture 链式编排即可,代码简洁直观
  • 步骤多、有分支/重试/回退逻辑:用状态机模式,或引入 spring statemachine 框架
  • 对实时性要求高:优先使用回调/webhook 模式,避免轮询带来的延迟
  • 对可靠性要求高:无论哪种方案,都要加上超时机制、失败重试、任务持久化(防止服务重启丢失流程状态)

工业现场的注意事项

在实际工业项目中,正常流程往往不是最难的,异常处理才是重中之重:

  • 超时处理:小车长时间未到达指定位置,要能自动告警并取消任务
  • 失败重试:夹取失败时要能自动重试(设置最大重试次数)
  • 断点恢复:服务重启后能恢复未完成的流程,而不是从头开始
  • 通信容错:网络断开后恢复时,要能继续未完成的流程
  • 日志追踪:每个步骤都要有完整的日志记录,方便排查问题

总结

java 实现 agv 多步骤异步任务编排,核心思路就是异步等待 + 状态驱动completablefuture 提供了轻量级的链式编排能力,状态机模式提供了复杂流程的可控性,而事件驱动则提供了最佳的实时性。根据实际业务复杂度选择合适的方案,同时务必做好异常处理和可靠性保障,这才是工业级应用的关键所在。

以上就是java实现多步骤异步任务编排的三种主流方案的详细内容,更多关于java多步骤异步任务编排的资料请关注代码网其它相关文章!

(0)

相关文章:

版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。

发表评论

验证码:
Copyright © 2017-2026  代码网 保留所有权利. 粤ICP备2024248653号
站长QQ:2386932994 | 联系邮箱:2386932994@qq.com