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

Spring+Ribbon环境下如何拦截请求用于数据库回滚重放?

针对Ribbon调用场景的请求拦截方案

刚好遇到过类似的多版本数据库兼容测试场景,给你几个针对性的方案,都是Spring/Ribbon生态里的标准玩法,应该能完美解决你的拦截PUT/POST/DELETE请求并存入Kafka的需求:

方案1:使用Ribbon原生RequestInterceptor(最推荐)

Ribbon本身提供了请求拦截机制,不管你是用RestTemplate配合@LoadBalanced,还是用Feign(底层依赖Ribbon),这个拦截器都能生效,而且只针对指定的Ribbon客户端(比如你的<D_NEW>),不会影响其他服务调用。

实现步骤:

  1. 实现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 省略
    }
}
  1. 绑定到指定的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;
    }
}

注意事项:

  1. 异步发送Kafka:一定要用异步方式发送,避免阻塞主请求的处理流程,影响接口性能。
  2. 请求体序列化:如果请求体是二进制格式(比如文件上传),需要特殊处理,不要强行转JSON。
  3. 异常隔离:拦截器中的异常要自行处理,不要抛出到主业务流程中,避免因为Kafka故障导致正常请求失败。
  4. 幂等性:后续重放请求到<D_OLD>时,要确保请求是幂等的,避免重复执行导致数据不一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 06:43:01