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

Spring Cloud Gateway后置过滤器封装微服务响应至ConsumerResponse

Spring Cloud Gateway后置过滤器包装微服务响应为ConsumerResponse

需求

需要在Spring Cloud Gateway的Post Filter中获取微服务返回的响应,将该响应数据放入自定义的ConsumerResponse对象中,作为最终返回给UI的响应结构。

微服务返回响应示例

{
    "userDetails": {
        "userId": 24,
        "userName": "ABC",
        "description": "ABC from ZZZ",
        "registrationStatus": "REGISTERED",
        "registrationTime": [
            2022,
            5,
            23,
            21,
            17,
            41,
            465000000
        ],
        "lastUpdatedTime": [
            2022,
            5,
            23,
            21,
            17,
            41,
            465000000
        ]
    }
}

ConsumerResponse类定义

@JsonInclude(JsonInclude.Include.NON_NULL)
@Data
public class ConsumerResponse implements Serializable {

  private static final long serialVersionUID = 1L;
  
  private ConsumerResponse() {}  

  public static <T> ResponseEntity<ResponseEnvelope<Object>> ok(
      T data, String apiVersion, Status status, HttpStatus httpStatus) {
    return ResponseEntity.status(httpStatus)
        .body(ResponseEnvelope.builder().apiVersion(apiVersion).data(data).status(status).build());
  }
}

现有Gateway Filter代码(待修改)

package org.xyz.gateway.filter;

import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.stream.Collectors;

import org.reactivestreams.Publisher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.gateway.filter.GatewayFilterChain;
import org.springframework.cloud.gateway.filter.GlobalFilter;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferFactory;
import org.springframework.core.io.buffer.DefaultDataBuffer;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.http.server.reactive.ServerHttpResponseDecorator;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ResponseStatusException;
import org.springframework.web.server.ServerWebExchange;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

@Component
public class AuthFilter implements GlobalFilter {
    
    private final Logger logger = LoggerFactory.getLogger(AuthFilter.class);
    
    @Autowired
    private IdService idService;
    
    @Override
    public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {    
                
                // code to fetch id from request
                
                MyResult response = idService.getEntitlements(id);              
                
                ArrayList<String> rolesList = response.getRoles()
                        .values()
                        .stream()
                        .distinct()
                        .collect(Collectors.toCollection(ArrayList::new));
        
                
                ServerHttpRequest httpRequest = exchange.getRequest()
                        .mutate()
                        .header("id", id)
                        .header("roles", rolesList.toString())
                        .build();
                
                ServerWebExchange mutatedExchange = exchange
                            .mutate()
                            .request(httpRequest)
                            .build();
                
            return chain.filter(mutatedExchange).then(Mono.fromRunnable(() -> {
                ServerHttpResponse serverHttpResponse = exchange.getResponse();
                HttpStatus responseStatus = serverHttpResponse.getStatusCode();              
                DataBuffer dataBuffer = null;
                String path = exchange.getRequest().getPath().toString();               
                ServerHttpRequest serverHttpRequest = exchange.getRequest();
                DataBufferFactory dataBufferFactory = serverHttpResponse.bufferFactory();
                
                ServerHttpResponseDecorator decoratedResponse = getDecoratedResponse(path, serverHttpResponse, serverHttpRequest, dataBufferFactory);
                try {
                    dataBuffer = dataBufferFactory.wrap(
                          new ObjectMapper().writeValueAsBytes(
                                  // need to set the response data here in ConsumerResponse Object below
                                  ConsumerResponse.ok(decoratedResponse.toString(), GatewayConstants.VERSION, Status.OK, responseStatus)));
                    } catch (JsonProcessingException e) {
                      dataBuffer = serverHttpResponse.bufferFactory().wrap("".getBytes());
                    }
                logger.info("decoratedResponse = {}", decoratedResponse);
                exchange.getResponse().getHeaders().setContentType(MediaType.APPLICATION_JSON);
                serverHttpResponse.getHeaders().setContentLength(decoratedResponse.toString().length());
                serverHttpResponse.writeWith(Mono.just(dataBuffer)).subscribe();
                exchange.mutate().response(serverHttpResponse).build();
            }));
    }
    
    private ServerHttpResponseDecorator getDecoratedResponse(String path, ServerHttpResponse response, ServerHttpRequest request, DataBufferFactory dataBufferFactory) {
        return new ServerHttpResponseDecorator(response) {

            @Override
            public Mono<Void> writeWith(final Publisher<? extends DataBuffer> body) {
                logger.info("writeWith entering(...)---> with body = {} ", body);
                if (body instanceof Flux) {
                    logger.info("Body is Flux");
                    Flux<? extends DataBuffer> fluxBody = (Flux<? extends DataBuffer>) body;

                    return super.writeWith(fluxBody.buffer().map(dataBuffers -> {

                        DefaultDataBuffer joinedBuffers = new DefaultDataBufferFactory().join(dataBuffers);
                        byte[] content = new byte[joinedBuffers.readableByteCount()];
                        joinedBuffers.read(content);
                         String responseBody = new String(content, StandardCharsets.UTF_8);//MODIFY RESPONSE and Return the Modified response
                        logger.debug("requestId: {}, method: {}, url: {}, \nresponse body :{}", request.getId(), request.getMethodValue(), request.getURI(), responseBody);

                        return dataBufferFactory.wrap(responseBody.getBytes());
                    })).onErrorResume(err -> {

                        logger.error("error while decorating Response: {}",err.getMessage());
                        return Mono.empty();
                    });

                }

                return super.writeWith(body);
            }
        };
    }
}

解决方案代码修改

原代码错误地尝试通过ServerHttpResponseDecorator.toString()获取响应体,且后置处理逻辑位置不合理。正确做法是在ServerHttpResponseDecorator的writeWith方法中直接处理微服务返回的响应体,包装为ConsumerResponse结构后返回。

修改后的完整AuthFilter代码:

package org.xyz.gateway.filter;

import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Map;
import java.util.stream.Collectors;

import org.reactivestreams.Publisher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.gateway.filter.GatewayFilterChain;
import org.springframework.cloud.gateway.filter.GlobalFilter;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferFactory;
import org.springframework.core.io.buffer.DefaultDataBuffer;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.http.server.reactive.ServerHttpResponseDecorator;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ServerWebExchange;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

@Component
public class AuthFilter implements GlobalFilter {
    
    private final Logger logger = LoggerFactory.getLogger(AuthFilter.class);
    private final ObjectMapper objectMapper = new ObjectMapper();
    
    @Autowired
    private IdService idService;
    
    @Override
    public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {    
        // 1. 从请求中获取id(替换为你的实际id获取逻辑)
        String id = "";
                
        // 2. 获取权限信息并设置请求头
        MyResult response = idService.getEntitlements(id);              
        ArrayList<String> rolesList = response.getRoles()
                .values()
                .stream()
                .distinct()
                .collect(Collectors.toCollection(ArrayList::new));
        
        ServerHttpRequest httpRequest = exchange.getRequest()
                .mutate()
                .header("id", id)
                .header("roles", rolesList.toString())
                .build();
        
        ServerWebExchange mutatedExchange = exchange
                    .mutate()
                    .request(httpRequest)
                    .build();
        
        // 3. 装饰响应,处理微服务返回结果并包装为ConsumerResponse结构
        ServerHttpResponse originalResponse = mutatedExchange.getResponse();
        DataBufferFactory bufferFactory = originalResponse.bufferFactory();
        
        ServerHttpResponseDecorator decoratedResponse = new ServerHttpResponseDecorator(originalResponse) {
            @Override
            public Mono<Void> writeWith(Publisher<? extends DataBuffer> body) {
                if (body instanceof Flux) {
                    Flux<? extends DataBuffer> fluxBody = (Flux<? extends DataBuffer>) body;
                    
                    return super.writeWith(fluxBody.buffer().map(dataBuffers -> {
                        // 读取微服务返回的响应体
                        DefaultDataBuffer joinedBuffers = new DefaultDataBufferFactory().join(dataBuffers);
                        byte[] content = new byte[joinedBuffers.readableByteCount()];
                        joinedBuffers.read(content);
                        String responseBodyStr = new String(content, StandardCharsets.UTF_8);
                        
                        logger.debug("微服务响应体: {}", responseBodyStr);
                        
                        try {
                            // 解析响应体为Object(如果有对应实体类,可替换为具体类型)
                            Object microserviceResponse = objectMapper.readValue(responseBodyStr, Map.class);
                            
                            // 构建ConsumerResponse对应的ResponseEnvelope
                            ResponseEnvelope<Object> envelope = ResponseEnvelope.builder()
                                    .apiVersion(GatewayConstants.VERSION)
                                    .data(microserviceResponse)
                                    .status(Status.OK)
                                    .build();
                            
                            // 将Envelope序列化为JSON字节数组
                            byte[] newResponseBytes = objectMapper.writeValueAsBytes(envelope);
                            
                            // 设置响应头
                            originalResponse.getHeaders().setContentType(MediaType.APPLICATION_JSON);
                            originalResponse.getHeaders().setContentLength(newResponseBytes.length);
                            
                            return bufferFactory.wrap(newResponseBytes);
                        } catch (JsonProcessingException e) {
                            logger.error("解析或序列化响应失败: {}", e.getMessage(), e);
                            // 异常场景返回空响应或自定义错误响应
                            return bufferFactory.wrap(new byte[0]);
                        }
                    })).onErrorResume(err -> {
                        logger.error("处理响应时发生异常: {}", err.getMessage(), err);
                        return Mono.empty();
                    });
                }
                return super.writeWith(body);
            }
        };
        
        // 4. 使用装饰后的响应继续执行过滤链
        return chain.filter(exchange.mutate().response(decoratedResponse).build());
    }
}

关键修改点说明

  • 移除冗余的getDecoratedResponse方法,直接在filter方法内创建ServerHttpResponseDecorator,逻辑更紧凑。
  • 在writeWith方法中直接读取微服务响应体字节,解析为Object(若有对应实体类,可替换为具体类型)。
  • 直接构建ResponseEnvelope替代ConsumerResponse.ok()(网关中无需构建ResponseEntity,直接操作响应体更适配WebFlux模型)。
  • 正确设置响应头的Content-Type和Content-Length,确保UI端能正确解析响应。
  • 增加异常处理逻辑,避免JSON解析失败导致网关报错。

内容的提问来源于stack exchange,提问作者VinodKhade

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 05:11:11