Spring Kafka @SendTo("!{someExpression}")运行时表达式的解释与示例咨询
解释@SendTo("!{someExpression}")运行时表达式及示例
核心逻辑说明
官方文档定义:
@SendTo("!{someExpression}")会在运行时通过计算表达式路由至对应主题。表达式计算的#root对象包含三个属性:
request:入站的ConsumerRecord(批量监听器场景为ConsumerRecords对象)source:由request转换而来的org.springframework.messaging.Message<?>对象result:当前方法的执行返回结果
简单来说,!{...}语法允许你不写死目标主题,而是根据消息请求、转换后的消息对象或者方法处理结果,动态计算出要发送的主题,极大提升路由灵活性。
实际使用示例
示例1:根据方法返回结果动态选择主题
假设业务中需要根据订单类型,将处理后的结果发送到不同主题。我们可以利用result对象的属性来计算目标主题:
@KafkaListener(topics = "order-input") @SendTo("!{result.targetTopic}") public OrderProcessingResult processOrder(ConsumerRecord<String, Order> record) { Order order = record.value(); OrderProcessingResult result = new OrderProcessingResult(); // 根据订单类型设置目标主题 if (order.isVip()) { result.setTargetTopic("vip-order-output"); } else { result.setTargetTopic("normal-order-output"); } // 其他处理逻辑... return result; }
方法执行后,Spring会计算result.targetTopic的值,将结果消息发送到对应的主题。
示例2:基于入站请求的元数据生成主题
如果需要将处理后的消息发送到原主题的衍生主题(比如原主题加-processed后缀),可以直接使用request对象的属性:
@KafkaListener(topics = "user-action-logs") @SendTo("!{request.topic + '-processed'}") public UserActionLog processUserLog(ConsumerRecord<String, UserActionLog> record) { UserActionLog log = record.getValue(); log.setProcessedTimestamp(System.currentTimeMillis()); // 日志处理逻辑... return log; }
当监听到user-action-logs主题的消息时,处理后的结果会自动发送到user-action-logs-processed主题。
示例3:通过消息头指定目标主题
如果入站消息的Header中已经携带了目标主题信息,可以通过source对象获取Header值来路由:
@KafkaListener(topics = "cross-service-events") @SendTo("!{source.headers['destination-topic']}") public ServiceEvent processCrossServiceEvent(Message<ServiceEvent> message) { ServiceEvent event = message.getPayload(); event.setHandled(true); // 跨服务事件处理逻辑... return event; }
这里source是转换后的Message对象,表达式直接读取Header中的destination-topic值,将结果发送到指定主题。
内容的提问来源于stack exchange,提问作者Amir Choubani
相关产品推荐
相关产品推荐

