Quarkus环境下Camel路由拆分与EIP选型技术咨询
Apache Camel + Quarkus 路由拆分与EIP选型问题解答
业务场景概述
- 从AMQ队列接收消息,重新映射后发送至不同客户的REST端点
- 每个客户配置包含字段筛选规则、OAuth地址、消息发送地址及凭证
- 客户按Agent分组,配置以
agentId为键、customerConfigs列表为值存储在Map中 - 根据消息中的字段确定所属Agent,遍历该Agent下所有客户,按客户需求重映射消息
- 按客户规则过滤消息,符合条件则完成OAuth认证后发送,否则跳过
当前面临的问题
- 应选用Recipient List还是Dynamic Router等EIP组件?
- 拆分路由时,RemappedMessage与CustomerConfig(一对多关系)如何传递至下游?
- 如何从Exchange对象中获取属性,动态配置
to()的URL及凭证?
解决方案
1. EIP组件选型:优先使用Recipient List
Recipient List更适配你的场景:
- 你的业务逻辑是明确根据Agent配置获取所有目标客户列表,属于"已知收件人集合"的场景
- Recipient List可直接基于预先生成的目标列表,将消息分发到多个端点,且支持对每个收件人单独处理
- Dynamic Router更适用于需要**动态决定下一个路由节点(不确定最终节点数量或需中途调整)**的场景,你的需求不需要这种动态决策能力,用Recipient List更简洁直接
2. 一对多对象的传递与路由拆分
通过Split组件+Exchange属性存储处理一对多关系:
- 在
CustomerConfigRetrieverBean中,将获取到的customerConfigs列表存入Exchange属性:exchange.setProperty("customerConfigs", customerConfigsList) - 在
EndpointFieldsTailor.class之后,用Split组件拆分customerConfigs列表,同时保留原消息(RemappedMessage):.split(simple("${exchangeProperty.customerConfigs}")) .shareUnitOfWork() // 可选,确保拆分后的分支共用同一事务上下文 .setProperty("currentCustomerConfig", simple("${body}")) // 将当前拆分出的CustomerConfig存入属性 .bean(MessageFilter.class) // 按客户规则过滤消息,不符合则终止当前分支 .bean(OAuthTokenRetriever.class) // 根据currentCustomerConfig获取OAuth token并存入属性 - 也可以封装
CustomerMessageWrapperDTO,包含RemappedMessage和CustomerConfig,将Wrapper列表作为Split的数据源,每个拆分后的Exchange的body即为单个Wrapper对象
3. 动态配置to()的URL及凭证
动态URL配置
使用Camel的toD()(动态to)方法,直接从Exchange属性读取客户的发送地址:
.toD("${exchangeProperty.currentCustomerConfig.sendUrl}")
动态OAuth凭证处理
- 在
OAuthTokenRetrieverBean中,根据currentCustomerConfig里的OAuth地址、clientId、clientSecret调用接口获取token,存入Exchange属性:exchange.setProperty("oauthToken", token) - 发送请求前设置请求头:
.setHeader("Authorization", simple("Bearer ${exchangeProperty.oauthToken}")) - 若需更灵活的OAuth集成,可使用Camel的
oauth2组件动态配置参数:.to("oauth2:${exchangeProperty.currentCustomerConfig.oauthUrl}?clientId=${exchangeProperty.currentCustomerConfig.clientId}&clientSecret=${exchangeProperty.currentCustomerConfig.clientSecret}")
完整路由示例(修改后)
from("activemq:queue:" + appConfig.getQueueName()) .bean(IncomingMessageConverter.class) .bean(UserIdValidator.class) // 验证不通过则终止路由 .bean(CustomerConfigRetrieverBean.class) // 将customerConfigs列表存入exchangeProperty .bean(EndpointFieldsTailor.class) // 生成RemappedMessage,保留在body中 .split(simple("${exchangeProperty.customerConfigs}")) .shareUnitOfWork() .setProperty("currentCustomerConfig", simple("${body}")) .bean(MessageFilter.class) // 过滤不符合规则的消息,当前分支终止 .bean(OAuthTokenRetriever.class) // 获取OAuth token并存入exchangeProperty .setHeader("Authorization", simple("Bearer ${exchangeProperty.oauthToken}")) .toD("${exchangeProperty.currentCustomerConfig.sendUrl}") // 配置仅重试发送步骤 .onException(Exception.class) .maximumRedeliveries(3) .redeliveryDelay(1000) .retryAttemptedLogLevel(LoggingLevel.WARN) .end()
重试策略说明
上述路由通过onException配置仅对发送阶段的异常进行重试,Split后的分支仅处理单个客户的发送逻辑,异常只会触发当前分支的重试,不会影响其他客户的消息处理。
内容的提问来源于stack exchange,提问作者WesternGun
相关产品推荐
相关产品推荐

