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
相关产品推荐
相关产品推荐

