You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 18:45:41