You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Apache Camel实现同步顺序路由处理及路由单次执行控制?

解决Camel路由顺序执行(依赖处理)问题

核心需求

通过YAML配置多个独立路由,支持:

  • 按匹配规则(文件/正则)处理文件
  • 执行移动、复制、压缩、删除、SFTP传输等操作
  • 严格顺序执行:必须等前一个路由完成所有文件处理后,再启动下一个路由;全部完成后循环重复执行

你当前方法的问题

你写的waitForRouteCompletion存在两个核心缺陷:

  1. 硬编码Thread.sleep(3000)完全不可靠,无法适配不同处理时长的场景
  2. 手动调用stopRoute没有配合文件组件的停止逻辑,路由仍可能响应已轮询到的文件

同时你尝试的start/continue/end、单线程、controlBus等方案不适用的原因是:这些机制都是针对单路由内的流程控制,或者路由的启停调度,没有解决路由间依赖的顺序触发问题。

正确实现方案

1. 核心配置:文件路由自动停止

给每个文件消费路由添加consumer.stopWhenEmpty=true配置,这样路由会在目录中没有匹配文件时自动停止,无需手动干预。

2. 自定义RoutePolicy控制路由启停顺序

创建一个路由策略,监听路由停止事件,自动启动下一个路由,最后一个路由停止后回到第一个,形成循环。

自定义RoutePolicy代码

public class SequentialRoutePolicy extends RoutePolicySupport {
    private List<String> routeExecutionOrder;
    private int currentRouteIndex = 0;

    @Override
    public void onRouteStopped(Route route) {
        super.onRouteStopped(route);
        // 切换到下一个路由索引,循环回到开头
        currentRouteIndex = (currentRouteIndex + 1) % routeExecutionOrder.size();
        try {
            getCamelContext().getRouteController().startRoute(routeExecutionOrder.get(currentRouteIndex));
        } catch (Exception e) {
            log.error("Failed to start next route: {}", routeExecutionOrder.get(currentRouteIndex), e);
        }
    }

    // 注入路由执行顺序列表
    public void setRouteExecutionOrder(List<String> routeExecutionOrder) {
        this.routeExecutionOrder = routeExecutionOrder;
    }
}

YAML路由配置示例

camel:
  routes:
    # 路由1:处理txt文件并移动到processed目录
    - id: file-process-route
      uri: "file://./input?include=*.txt&delete=true&consumer.stopWhenEmpty=true"
      steps:
        - to: "file://./processed"
      routePolicy:
        ref: sequentialRoutePolicy

    # 路由2:压缩processed目录的txt文件
    - id: file-compress-route
      uri: "file://./processed?include=*.txt&compression=gzip&consumer.stopWhenEmpty=true"
      steps:
        - to: "file://./compressed"
      routePolicy:
        ref: sequentialRoutePolicy

    # 路由3:将压缩文件上传到SFTP
    - id: sftp-upload-route
      uri: "file://./compressed?include=*.gz&consumer.stopWhenEmpty=true"
      steps:
        - to: "sftp://your-user@your-host/remote-path?password=your-password"
      routePolicy:
        ref: sequentialRoutePolicy

  # 注册自定义路由策略
  registry:
    beans:
      sequentialRoutePolicy:
        type: com.your-package.SequentialRoutePolicy
        properties:
          routeExecutionOrder: ["file-process-route", "file-compress-route", "sftp-upload-route"]

3. 初始启动控制

在应用启动时,只启动第一个路由,其他路由保持停止状态:

public class CamelApp {
    public static void main(String[] args) throws Exception {
        CamelContext camelContext = new DefaultCamelContext();
        // 加载YAML路由配置(根据你的实际加载方式调整)
        YamlRoutesDefinitionLoader loader = new YamlRoutesDefinitionLoader(camelContext);
        loader.loadRoutes("classpath:routes.yaml");

        RouteController routeController = camelContext.getRouteController();
        // 启动第一个路由
        routeController.startRoute("file-process-route");
        // 停止其他路由(初始不启动)
        routeController.stopRoute("file-compress-route");
        routeController.stopRoute("sftp-upload-route");

        camelContext.start();
        // 保持应用运行
        Thread.currentThread().join();
    }
}

工作流程说明

  1. 启动后只有第一个路由运行,处理input目录下所有txt文件,处理完成后因stopWhenEmpty=true自动停止
  2. 路由停止事件触发SequentialRoutePolicy的onRouteStopped方法,自动启动第二个路由
  3. 第二个路由处理processed目录下的txt文件,压缩完成后自动停止,触发启动第三个路由
  4. 第三个路由上传压缩文件到SFTP,完成后自动停止,触发回到第一个路由,循环执行

额外注意事项

  • 如果需要路由处理完一次后不再循环,只需去掉currentRouteIndex = (currentRouteIndex + 1) % routeExecutionOrder.size()中的取模逻辑,改为判断是否到达最后一个路由即可
  • 对于SFTP等远程组件,确保配置了正确的连接参数和错误重试机制
  • 可以给路由添加consumer.maxMessagesPerPoll参数,控制每次轮询处理的文件数量,避免一次性处理过多文件

内容的提问来源于stack exchange,提问作者SomeRandomGuy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 01:05:10