Spring中如何合并RabbitMQ消息与接口数据填充至同一DTO模型类
实现步骤
1 先修正PeopleDocumentDTO的字段映射与序列化问题
现有DTO存在3个核心问题:
- JSON字段名和类属性名不匹配,会导致反序列化丢值
- 内部类
Customer是非静态类,Jackson无法直接实例化 - 缺少对应嵌套
Document类的定义,无法匹配接收的JSON结构
修正后的DTO代码如下:
import com.fasterxml.jackson.annotation.JsonProperty; import lombok.Getter; import lombok.Setter; import java.util.List; @Getter @Setter public class PeopleDocumentDTO { @JsonProperty("type") private String processType; private String operation; private String entity; private String entityType; private Long id; @JsonProperty("documents") private Document document; @Getter @Setter public static class Customer { private String systemId; private String customerId; } // 补全缺失的Document内部类结构,匹配接收的JSON格式 @Getter @Setter public static class Document { private Long id; private Additionals additionals; private String code; private String typeDocument; } @Getter @Setter public static class Additionals { private String issuing_authority; private String country_doc; private String place_of_birth; private String valid_from; private String valid_to; } @JsonProperty("relatedCustomers") private List<Customer> customers; }
2 调整RabbitListener消费逻辑,填充关联客户数据
WebClient是响应式客户端,在RabbitMQ同步消费场景下需要通过block()方法同步获取返回结果,同时要做实体类型转换,把接口返回的CustomerRelation转成DTO要求的Customer结构。
优化点:
- 不要每次消费消息都新建ObjectMapper,直接注入Spring容器托管的单例ObjectMapper
- 调用getPerson拿到过滤后的关联客户后,遍历转换类型赋值给DTO
- 授权Token从配置文件读取,避免硬编码
修正后的消费方法代码:
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import java.nio.charset.StandardCharsets; import java.util.List; import java.util.stream.Collectors; @Component public class DocumentMessageListener { private final ObjectMapper objectMapper; private final WebClientService webClientService; // 从配置文件读取外部接口授权token @Value("${external.api.auth-token}") private String GS_AUTH_TOKEN; // 构造注入依赖 public DocumentMessageListener(ObjectMapper objectMapper, WebClientService webClientService) { this.objectMapper = objectMapper; this.webClientService = webClientService; } @RabbitListener(queues = "${event.queue}") public void receivedMessage(Message message) throws Exception { // 1. 解析消息体为DTO String json = new String(message.getBody(), StandardCharsets.UTF_8); PeopleDocumentDTO dto = objectMapper.readValue(json, PeopleDocumentDTO.class); // 非法消息直接拦截,避免空指针 if (dto.getId() == null) { return; } // 2. 调用外部接口获取过滤后的关联客户 Person person = webClientService.getPerson(dto.getId().intValue(), GS_AUTH_TOKEN) .block(); // 同步阻塞等待响应式结果返回,适配MQ同步消费场景 if (person != null && person.getRelatedCustomers() != null) { // 3. 将CustomerRelation转换为DTO要求的Customer类型,填充到DTO List<PeopleDocumentDTO.Customer> customerList = person.getRelatedCustomers() .stream() .map(relation -> { PeopleDocumentDTO.Customer customer = new PeopleDocumentDTO.Customer(); // 类型转换:relation中systemId是Integer,Customer中定义为String customer.setSystemId(String.valueOf(relation.getSystemId())); customer.setCustomerId(relation.getCustomerId()); return customer; }) .collect(Collectors.toList()); dto.setCustomers(customerList); } // 后续写DTO的业务处理逻辑即可,此处可打印最终结构验证 System.out.println(objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(dto)); } }
关键注意点
block()方法禁止在响应式异步调度线程中调用,在RabbitListener的同步消费线程中使用是安全的- 如果配置ObjectMapper开启下划线转驼峰命名策略,可以省略大部分@JsonProperty注解
- 你给出的最终JSON示例中重复出现的
id字段属于示例笔误,实际序列化时不会出现重复字段 - 可以根据业务需要给WebClient调用加超时、异常重试、失败降级逻辑,避免外部接口异常影响MQ消费
内容的提问来源于stack exchange,提问作者DiegoMG
相关产品推荐
相关产品推荐

