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

Spring Integration缓存实现求助:避免重复调用Rest API

实现Spring Integration中的API结果缓存

搞定这个缓存需求其实不难,我给你一步步拆解实现方案:


1. 先确保你的Payload类能作为HashMap的键

HashMap是通过equals()和hashCode()方法判断键是否重复的,所以你的自定义Payload类必须正确重写这两个方法,基于业务上唯一标识Payload的字段来实现,不然即使内容完全相同的两个Payload对象,也会被HashMap当成不同的键。

比如你的Payload类可以这么写:

import java.util.Objects;

public class Payload {
    // 假设这些是Payload的核心属性
    private String requestId;
    private String data;

    // Getter、Setter省略

    @Override
    public boolean equals(Object o) {
        if (this == o) return true;
        if (o == null || getClass() != o.getClass()) return false;
        Payload payload = (Payload) o;
        // 用核心属性判断相等
        return Objects.equals(requestId, payload.requestId) &&
               Objects.equals(data, payload.data);
    }

    @Override
    public int hashCode() {
        // 基于相同的核心属性生成哈希值
        return Objects.hash(requestId, data);
    }
}

2. 创建线程安全的缓存容器

因为Spring Integration的消息流通常是多线程处理的,普通HashMap在并发场景下会有线程安全问题,所以我们用ConcurrentHashMap来做缓存,把它注册成Spring Bean方便注入:

import java.util.concurrent.ConcurrentHashMap;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class CacheConfig {
    // 泛型替换成你的Payload类型和API返回结果类型
    @Bean
    public ConcurrentHashMap<Payload, ApiResponse> apiResponseCache() {
        return new ConcurrentHashMap<>();
    }
}

3. 修改IntegrationFlow,加入缓存逻辑

在调用Rest API之前,先检查缓存:如果缓存里有对应Payload的结果,直接返回;没有就调用API,然后把结果存入缓存。我们可以用handle()方法来封装这个逻辑:

import org.springframework.context.annotation.Bean;
import org.springframework.integration.dsl.IntegrationFlows;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.messaging.MessageHeaders;
import org.springframework.http.HttpHeaders;
import org.springframework.http.MediaType;
import org.springframework.web.client.RestTemplate;
import java.util.concurrent.ConcurrentHashMap;

@Bean
public IntegrationFlow read(ConcurrentHashMap<Payload, ApiResponse> apiResponseCache) {
    return IntegrationFlows.from("input")
            .split(new PayLoadSplitter())
            .enrichHeaders(p -> p.header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE))
            // 新增缓存检查与API调用逻辑
            .handle((payload, headers) -> {
                Payload requestPayload = (Payload) payload;
                // 第一步:查缓存
                ApiResponse cachedResult = apiResponseCache.get(requestPayload);
                if (cachedResult != null) {
                    return cachedResult;
                }
                // 第二步:缓存无数据,调用Rest API
                ApiResponse apiResult = callRestApi(requestPayload, headers);
                // 第三步:将结果存入缓存
                apiResponseCache.put(requestPayload, apiResult);
                return apiResult;
            })
            // 这里继续你的后续处理逻辑
            // .xxx()
            .get();
}

// 把原来的API调用逻辑抽成单独方法,便于维护
private ApiResponse callRestApi(Payload payload, MessageHeaders headers) {
    // 这里替换成你实际的Rest API调用代码,比如用RestTemplate或WebClient
    RestTemplate restTemplate = new RestTemplate();
    org.springframework.http.HttpHeaders requestHeaders = new org.springframework.http.HttpHeaders();
    requestHeaders.set(HttpHeaders.ACCEPT, headers.get(HttpHeaders.ACCEPT).toString());
    org.springframework.http.HttpEntity<Payload> requestEntity = new org.springframework.http.HttpEntity<>(payload, requestHeaders);
    
    return restTemplate.postForObject("https://your-api-endpoint-url", requestEntity, ApiResponse.class);
}

额外注意事项

  • 缓存过期策略:如果你的API数据会更新,单纯的HashMap不会自动清理过期数据。可以考虑:
    • 用Guava Cache或Caffeine Cache替代ConcurrentHashMap,它们支持自动过期配置
    • 自己给缓存值包装一层过期时间,每次查缓存时判断是否过期,过期则重新调用API并更新缓存
  • 缓存大小限制:如果Payload数量很大,HashMap可能会占用过多内存,可以给ConcurrentHashMap设置初始容量,或者定期清理旧数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:32:24