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

如何在源无消息时强制Aggregator按时间间隔执行?

问题1:修改现有方案满足无文件时发送0的需求

你的现有路由核心问题是聚合器在周期内无消息时不会触发完成逻辑,导致overviewProcessor不会执行。可以通过Camel聚合器的sendEmptyMessageWhenIdle()配置解决,同时调整聚合策略处理空消息场景:

1. 修改路由配置

在聚合器的配置中添加.sendEmptyMessageWhenIdle(),开启空闲时发送空消息的机制:

from(getSourceUri())
        .onException(Throwable.class).bean(exceptionHandler).handled(true).end()
        .process(sftpProcessor)
        .aggregate(constant(true), new OverviewAggregatorStrategy())
            .aggregationStrategyMethodAllowNull()
            .completionInterval(TWO_MINUTES)
            .sendEmptyMessageWhenIdle() // 新增:周期内无消息时发送空消息触发聚合
        .process(overviewProcessor)
        .end();

2. 调整聚合策略逻辑

更新OverviewAggregatorStrategy,处理空消息场景(即newExchange为null的情况),返回包含0计数的Exchange:

public class OverviewAggregatorStrategy implements AggregationStrategy {
    private static final String PROCESSED_MESSAGES = "PROCESSED_MESSAGES";

    @Override
    public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
        // 处理空闲时的空消息
        if (newExchange == null) {
            Exchange emptyExchange = oldExchange != null 
                ? new DefaultExchange(oldExchange.getContext()) 
                : new DefaultExchange(new SimpleCamelContext());
            emptyExchange.setProperty(PROCESSED_MESSAGES, 0L);
            return emptyExchange;
        }

        // 原有消息计数逻辑
        Long processedMessages;
        if (oldExchange != null) {
            processedMessages = oldExchange.getProperty(PROCESSED_MESSAGES, Long.class);
            oldExchange.setProperty(PROCESSED_MESSAGES, processedMessages + 1);
            return oldExchange;
        } else {
            newExchange.setProperty(PROCESSED_MESSAGES, 1L);
            return newExchange;
        }
    }
}

这样配置后,即使周期内没有SFTP文件消息,聚合器也会触发一次空聚合,overviewProcessor会收到计数为0的Exchange。


问题2:实现此类需求的更优方案

用聚合器依赖消息触发的方式绕了弯路,更直接的方案是定时主动统计,逻辑更清晰且易于维护,推荐两种实现思路:

思路1:定时扫描SFTP源统计文件数

直接用Camel的timer组件定期触发SFTP文件扫描,统计数量后发送到监控工具:

from("timer:sftpStats?period=" + TWO_MINUTES)
        .onException(Throwable.class).bean(exceptionHandler).handled(true).end()
        .process(exchange -> {
            // 调用SFTP客户端获取文件列表并统计数量
            // 示例:使用Camel SFTP组件的list操作
            List<Map<String, Object>> fileList = exchange.getContext()
                .createProducerTemplate()
                .requestBody(getSourceUri() + "?noop=true", null, List.class);
            long fileCount = fileList != null ? fileList.size() : 0;
            exchange.setProperty("PROCESSED_MESSAGES", fileCount);
        })
        .process(overviewProcessor)
        .end();

优点:逻辑直接,无需依赖消息触发,无论有无文件都会定时输出结果,避免聚合器的复杂配置。

思路2:分离文件处理与计数统计

如果需要统计已处理完成的文件数(而非仅存在的文件数),可以拆分路由:

  1. 文件处理路由:处理SFTP文件时递增计数器
  2. 统计路由:定时读取计数器值并重置,发送到监控工具
// 1. 文件处理路由:处理文件并计数
@Component
public class SftpProcessingRoute extends RouteBuilder {
    private final CounterService counterService;

    public SftpProcessingRoute(CounterService counterService) {
        this.counterService = counterService;
    }

    @Override
    public void configure() throws Exception {
        from(getSourceUri())
                .onException(Throwable.class).bean(exceptionHandler).handled(true).end()
                // 幂等消费避免重复处理
                .idempotentConsumer(header("CamelFileName"), 
                    FileIdempotentRepository.fileIdempotentRepository(new File(".processed-files")))
                .process(sftpProcessor)
                .bean(counterService, "increment")
                .end();
    }
}

// 2. 定时统计路由
@Component
public class StatsReportRoute extends RouteBuilder {
    private final CounterService counterService;
    private final OverviewProcessor overviewProcessor;

    public StatsReportRoute(CounterService counterService, OverviewProcessor overviewProcessor) {
        this.counterService = counterService;
        this.overviewProcessor = overviewProcessor;
    }

    @Override
    public void configure() throws Exception {
        from("timer:statsReport?period=" + TWO_MINUTES)
                .process(exchange -> {
                    // 获取当前计数并重置为0
                    long processedCount = counterService.getAndReset();
                    exchange.setProperty("PROCESSED_MESSAGES", processedCount);
                })
                .process(overviewProcessor)
                .end();
    }
}

// 计数器服务(单例)
@Component
public class CounterService {
    private final AtomicLong count = new AtomicLong(0);

    public void increment() {
        count.incrementAndGet();
    }

    public long getAndReset() {
        return count.getAndSet(0);
    }
}

优点:职责分离,文件处理和统计逻辑解耦,计数准确且支持幂等,适合需要跟踪已处理文件的场景。


内容的提问来源于stack exchange,提问作者Rafał Trójczak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 12:03:09