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

