【问题标题】:Apache Camel aggregation completion not workingApache Camel 聚合完成不起作用
【发布时间】:2022-10-13 00:42:20
【问题描述】:

我已经配置了一个路由来从交易所提取一些数据并聚合它们;这里是简单的总结:

@Component
@RequiredArgsConstructor
public class FingerprintHistoryRouteBuilder extends RouteBuilder {

    private final FingerprintHistoryService fingerprintHistoryService;

    @Override
    public void configure() throws Exception {
        from(FINGERPRINT_HISTORY_ENDPOINT)
                .aggregate( (AggregationStrategy) (oldExchange, newExchange) -> {
                    final FingerprintHistory newFingerprint = extract(newExchange);
                    if (oldExchange == null) {
                        List<FingerprintHistory> fingerprintHistories = new ArrayList<>();
                        fingerprintHistories.add(newFingerprint);
                        newExchange.getMessage().setBody(fingerprintHistories);
                        return newExchange;
                    }

                    final Message oldMessage = oldExchange.getMessage();
                    final List<FingerprintHistory> fingerprintHistories = (List<FingerprintHistory>) oldMessage.getBody(List.class);
                    fingerprintHistories.add(newFingerprint);

                    return oldExchange;
                })
                .constant(true)
                .completionSize(aggregateCount)
                .completionInterval(aggregateDuration.toMillis())
                .to(FINGERPRINT_PROCESS_AGGREGATION)
                .end();

        from(FINGERPRINT_PROCESS_AGGREGATION)
                .process(exchange -> {
                    List<FingerprintHistory> fingerprintHistories = exchange.getMessage().getBody(List.class);
                    fingerprintHistoryService.saveAll(fingerprintHistories);
                });
strong text
    }

}

问题是聚合完成永远不会起作用,例如这是我的测试示例:

@SpringBootTest
class FingerprintHistoryRouteBuilderTest {

    @Autowired
    ProducerTemplate producerTemplate;

    @Autowired
    FingerprintHistoryRouteBuilder fingerprintHistoryRouteBuilder;

    @Autowired
    CamelContext camelContext;

    @MockBean
    FingerprintHistoryService historyService;

    @Test
    void api_whenAggregate() {
        UserSearchActivity activity = ActivityFactory.buildSampleSearchActivity("127.0.0.1", "salam", "finger");
        Exchange exchange = buildExchange();
        exchange.getMessage().setBody(activity);

ReflectionTestUtils.setField(fingerprintHistoryRouteBuilder, "aggregateCount", 1); ReflectionTestUtils.setFiled(fingerprintHistoryRouteBuilder, "aggregateDuration", Duration.ofNanos(1)); producerTemplate.send(FingerprintHistoryRouteBuilder.FINGERPRINT_HISTORY_ENDPOINT, 交换); Mockito.verify(historyService).saveAll(Mockito.any()); }

    Exchange buildExchange() {
        DefaultExchange defaultExchange = new DefaultExchange(camelContext);
        defaultExchange.setMessage(new DefaultMessage(camelContext));
        return defaultExchange;
    }

}

结果如下:

需要但未调用:fingerprintHistoryService bean.saveAll( );

【问题讨论】:

    标签: java spring-boot apache-camel aggregate


    【解决方案1】:

    我构建了这个simplified example,并且测试通过了,所以看起来您对聚合的使用可能是正确的。

    您是否考虑过您的Mockito.verify() 呼叫是在交换完成路由之前发生的?您可以通过删除验证调用并将.log() 语句添加到 FINGERPRINT_PROCESS_AGGREGATION 路由来测试这一点。如果您在执行过程中看到日志输出,则您知道交换正在按预期进行路由。如果是这种情况,那么您的verify() 呼叫需要能够等待交换完成路由。我不怎么使用mockito,但看起来你可以这样做:

    Mockito.verify(historyService, timeout(10000)).saveAll(Mockito.any());
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-07-19
      • 1970-01-01
      • 2014-10-10
      • 2023-03-26
      • 2018-01-04
      相关资源
      最近更新 更多