Camel如何实现direct端点的真正同步,解决aggregate异步执行问题
Camel Direct路由搭配Aggregate组件消息无序问题
环境信息
- 框架版本:Camel 3.11 ~ 3.12.0-SNAPSHOT
- JDK版本:Java 16
需求说明
本示例中timer会发送两条消息,当前第一条消息尚未完成处理时,第二条消息就已开始执行,期望单条DIRECT路由中的消息可以按顺序逐个处理。
代码实现
@Override public void configure() throws Exception { from("timer:cTimer?period=1&repeatCount=2") .process(e -> { final Integer oldBody = timerMessageId.getAndIncrement(); log.info("Send message: {}", oldBody); e.getIn().setHeader(KEY, oldBody); final List<Integer> list = IntStream.range(1, 10).boxed().collect(Collectors.toList()); e.getIn().setBody(list); }) .to("direct:direct-order-file-process-route:" + psId); from(direct("direct-order-file-process-route:" + psId)) .routeId("routeId:" + psId) .process(exchange -> { log.info("Started process message: {}", exchange.getIn().getHeader(KEY)); }) .split() .body() .aggregate(header(KEY), this::aggregatorExchanges) .completionSize(5) .completionTimeout(1000L) .process(exchange -> { Thread.sleep(100); log.info("DONE: {} / {}", exchange.getIn().getHeader(KEY), exchange.getIn().getBody()); }) .end(); } private Exchange aggregatorExchanges(Exchange oldExchange, Exchange newExchange) { List<Integer> drinks; if (oldExchange == null) { drinks = new ArrayList<>(); } else { drinks = (List<Integer>) oldExchange.getIn().getBody(); } drinks.add((Integer) newExchange.getIn().getBody()); newExchange.getIn().setBody(drinks); return newExchange; }
运行日志对比
实际运行日志
Send message: 1 Started process message: 1 Send message: 2 Started process message: 2 DONE: 1 / [1, 2, 3, 4, 5] DONE: 2 / [1, 2, 3, 4, 5] DONE: 1 / [6, 7, 8, 9] DONE: 2 / [6, 7, 8, 9]
期望运行日志
Send message: 1 Started process message: 1 DONE: 1 / [1, 2, 3, 4, 5] DONE: 1 / [6, 7, 8, 9] Send message: 2 Started process message: 2 DONE: 2 / [1, 2, 3, 4, 5] DONE: 2 / [6, 7, 8, 9]
问题原因
当前问题与aggregate组件相关,未使用aggregate时direct运行逻辑符合预期。换言之,direct与aggregate未同步(二者运行在不同线程中),direct不会等待聚合操作完成。
内容的提问来源于stack exchange,提问作者Sitnikov Artem
相关产品推荐
相关产品推荐

