Spring+Ribbon环境下如何拦截请求用于数据库回滚重放?
针对Ribbon调用场景的请求拦截方案
刚好遇到过类似的多版本数据库兼容测试场景,给你几个针对性的方案,都是Spring/Ribbon生态里的标准玩法,应该能完美解决你的拦截PUT/POST/DELETE请求并存入Kafka的需求:
方案1:使用Ribbon原生RequestInterceptor(最推荐)
Ribbon本身提供了请求拦截机制,不管你是用RestTemplate配合@LoadBalanced,还是用Feign(底层依赖Ribbon),这个拦截器都能生效,而且只针对指定的Ribbon客户端(比如你的<D_NEW>),不会影响其他服务调用。
实现步骤:
- 实现
RequestInterceptor接口,在apply方法中判断请求方法,序列化请求信息发送到Kafka:
@Component public class RibbonDataSyncInterceptor implements RequestInterceptor { private static final Logger log = LoggerFactory.getLogger(RibbonDataSyncInterceptor.class); private final KafkaTemplate<String, String> kafkaTemplate; private final ObjectMapper objectMapper; // 通过构造注入依赖 public RibbonDataSyncInterceptor(KafkaTemplate<String, String> kafkaTemplate, ObjectMapper objectMapper) { this.kafkaTemplate = kafkaTemplate; this.objectMapper = objectMapper; } @Override public void apply(RequestTemplate requestTemplate) { HttpMethod method = HttpMethod.valueOf(requestTemplate.getMethod()); // 只拦截写操作请求 if (Set.of(HttpMethod.POST, HttpMethod.PUT, HttpMethod.DELETE).contains(method)) { try { // 封装请求元数据和体内容 SyncRequestLog logDto = new SyncRequestLog(); logDto.setMethod(method.name()); logDto.setTargetUri(requestTemplate.getUri()); logDto.setHeaders(requestTemplate.getHeaders()); logDto.setRequestBody(requestTemplate.getBody() != null ? objectMapper.readTree(requestTemplate.getBody()) : null); // 异步发送到Kafka,避免阻塞主请求 CompletableFuture.runAsync(() -> kafkaTemplate.send("d_new_sync_topic", objectMapper.writeValueAsString(logDto)) ).exceptionally(ex -> { // 异常处理:记录日志而非抛出,避免影响主业务 log.error("Failed to send request to Kafka", ex); return null; }); } catch (JsonProcessingException e) { log.error("Failed to serialize request log", e); } } } // 内部类用于封装同步日志 private static class SyncRequestLog { private String method; private URI targetUri; private HttpHeaders headers; private JsonNode requestBody; // Getter & Setter 省略 } }
- 绑定到指定的Ribbon客户端:
创建一个Ribbon配置类,将拦截器注册为Bean,然后通过@RibbonClient指定给<D_NEW>服务:
@Configuration public class DNewRibbonConfig { @Bean public RequestInterceptor ribbonDataSyncInterceptor(KafkaTemplate<String, String> kafkaTemplate, ObjectMapper objectMapper) { return new RibbonDataSyncInterceptor(kafkaTemplate, objectMapper); } } // 在你的Spring启动类或配置类上添加注解,指定<D_NEW>的客户端配置 @RibbonClient(name = "D_NEW", configuration = DNewRibbonConfig.class)
方案2:使用Feign RequestInterceptor(如果用Feign调用<D_NEW>)
如果你的服务层是用FeignClient来调用<D_NEW>的,直接用Feign的拦截器会更贴合Feign的调用流程,逻辑和Ribbon拦截器类似:
@Component public class FeignDataSyncInterceptor implements feign.RequestInterceptor { private static final Logger log = LoggerFactory.getLogger(FeignDataSyncInterceptor.class); private final KafkaTemplate<String, String> kafkaTemplate; private final ObjectMapper objectMapper; public FeignDataSyncInterceptor(KafkaTemplate<String, String> kafkaTemplate, ObjectMapper objectMapper) { this.kafkaTemplate = kafkaTemplate; this.objectMapper = objectMapper; } @Override public void apply(RequestTemplate template) { String method = template.method(); if (Set.of("POST", "PUT", "DELETE").contains(method.toUpperCase())) { try { SyncRequestLog logDto = new SyncRequestLog(); logDto.setMethod(method); logDto.setTargetUri(URI.create(template.url())); // 转换Feign headers为Spring HttpHeaders HttpHeaders headers = new HttpHeaders(); template.headers().forEach((key, values) -> values.forEach(val -> headers.add(key, val))); logDto.setHeaders(headers); logDto.setRequestBody(template.body() != null ? objectMapper.readTree(template.body()) : null); // 异步发送到Kafka CompletableFuture.runAsync(() -> kafkaTemplate.send("d_new_sync_topic", objectMapper.writeValueAsString(logDto)) ).exceptionally(ex -> { log.error("Kafka send failed", ex); return null; }); } catch (JsonProcessingException e) { log.error("Request serialization failed", e); } } } // 同样的SyncRequestLog内部类 }
这个拦截器会自动作用于所有FeignClient,如果你只想针对<D_NEW>,可以在@FeignClient注解中指定configuration属性绑定这个拦截器。
方案3:Tomcat全局请求拦截(不推荐,除非需要全局拦截)
如果你想在Tomcat层面拦截所有对外请求,也可以用Tomcat的Valve机制,但这个是全局生效的,需要额外判断请求是否发往<D_NEW>,而且要处理请求体的重复读取问题(因为Tomcat的请求流只能读一次):
public class TomcatDataSyncValve extends ValveBase { private static final Logger log = LoggerFactory.getLogger(TomcatDataSyncValve.class); private final KafkaTemplate<String, String> kafkaTemplate; private final ObjectMapper objectMapper; // 替换为<D_NEW>的服务地址或标识 private static final String TARGET_SERVICE_ID = "d-new-service"; public TomcatDataSyncValve(KafkaTemplate<String, String> kafkaTemplate, ObjectMapper objectMapper) { this.kafkaTemplate = kafkaTemplate; this.objectMapper = objectMapper; } @Override public void invoke(Request request, Response response) throws IOException, ServletException { HttpServletRequest servletRequest = (HttpServletRequest) request; // 包装请求,支持重复读取请求体 ContentCachingRequestWrapper wrappedRequest = new ContentCachingRequestWrapper(servletRequest); // 判断是否是发往<D_NEW>的请求(根据Host或请求路径判断) String requestUrl = wrappedRequest.getRequestURL().toString(); if (requestUrl.contains(TARGET_SERVICE_ID)) { String method = wrappedRequest.getMethod(); if (Set.of("POST", "PUT", "DELETE").contains(method)) { try { SyncRequestLog logDto = new SyncRequestLog(); logDto.setMethod(method); logDto.setTargetUri(new URI(requestUrl)); // 读取请求头 HttpHeaders headers = new HttpHeaders(); Enumeration<String> headerNames = wrappedRequest.getHeaderNames(); while (headerNames.hasMoreElements()) { String headerName = headerNames.nextElement(); headers.add(headerName, wrappedRequest.getHeader(headerName)); } logDto.setHeaders(headers); // 读取缓存的请求体 byte[] bodyBytes = wrappedRequest.getContentAsByteArray(); logDto.setRequestBody(bodyBytes.length > 0 ? objectMapper.readTree(bodyBytes) : null); // 异步发送到Kafka CompletableFuture.runAsync(() -> kafkaTemplate.send("d_new_sync_topic", objectMapper.writeValueAsString(logDto)) ).exceptionally(ex -> { log.error("Kafka send failed", ex); return null; }); } catch (Exception e) { log.error("Request logging failed", e); } } } // 继续处理请求 getNext().invoke(wrappedRequest, response); } }
然后在Spring Boot中注册这个Valve:
@Configuration public class TomcatConfig { @Bean public TomcatServletWebServerFactory tomcatFactory(KafkaTemplate<String, String> kafkaTemplate, ObjectMapper objectMapper) { TomcatServletWebServerFactory factory = new TomcatServletWebServerFactory(); factory.addEngineValves(new TomcatDataSyncValve(kafkaTemplate, objectMapper)); return factory; } }
注意事项:
- 异步发送Kafka:一定要用异步方式发送,避免阻塞主请求的处理流程,影响接口性能。
- 请求体序列化:如果请求体是二进制格式(比如文件上传),需要特殊处理,不要强行转JSON。
- 异常隔离:拦截器中的异常要自行处理,不要抛出到主业务流程中,避免因为Kafka故障导致正常请求失败。
- 幂等性:后续重放请求到<D_OLD>时,要确保请求是幂等的,避免重复执行导致数据不一致。
内容的提问来源于stack exchange,提问作者Renjith
相关产品推荐
相关产品推荐

