Apache Camel中基于Caffeine Cache黑名单过滤Kafka用户消息的实现及优化问询
先帮你梳理下当前场景:你已经搭建了两条Camel路由,一条在应用启动时加载黑名单文件到Caffeine Cache(缓存key为blacklistedIds,值是List<String>类型),另一条从Kafka消费User消息,需要基于这个黑名单完成过滤。下面逐个解决你的疑问:
问题1:如何利用Caffeine Cache中的blacklistedIds对入站消息进行过滤?
有两种实用的实现方式,按需选择:
方式1:路由内直接通过Caffeine组件获取缓存并过滤
在Unmarshal之后,先从Caffeine Cache取出黑名单列表,再用filter断言判断当前消息的userId不在列表中即可。示例代码如下:
from(kafka(...)) .id(getRouteId()) .unmarshal().json(JsonLibrary.Jackson, UserMessage.class) // 从Caffeine Cache获取黑名单列表,存入Exchange属性 .setProperty("blacklistedIds") .to("caffeine-cache://blacklistedIds?action=GET&key=blacklistedIds") // 过滤掉userId在黑名单中的消息(注意id类型转换,确保和黑名单字符串匹配) .filter(simple("${body.id.toString()} not in ${exchangeProperty.blacklistedIds}")) .to(output())
方式2:通过你提到的UserMessageFilterService Bean过滤
这就是你注释掉的实现思路,只要在Bean方法里完成黑名单校验逻辑就行,后面的问题2会详细讲怎么在Bean里获取缓存。
问题2:如何在UserMessageFilterService#isAuthorizedUser Bean方法中获取该Caffeine Cache?
推荐两种常用方式,适配不同的项目环境:
方式1:Spring注入缓存实例(Spring/Spring Boot项目首选)
直接通过@Autowired注入CaffeineCacheManager或者指定名称的Cache实例,代码示例:
@Service public class UserMessageFilterService { // 方式A:注入CacheManager,按需获取指定缓存 private final CaffeineCacheManager cacheManager; public UserMessageFilterService(CaffeineCacheManager cacheManager) { this.cacheManager = cacheManager; } // 方式B:直接注入指定名称的缓存(需确保缓存已初始化) // @Autowired // private Cache blacklistedIdsCache; public boolean isAuthorizedUser(UserMessage userMessage) { // 方式A获取黑名单列表 Cache cache = cacheManager.getCache("blacklistedIds"); List<String> blacklistedIds = (List<String>) cache.get("blacklistedIds").get(); // 方式B直接获取 // List<String> blacklistedIds = (List<String>) blacklistedIdsCache.get("blacklistedIds").get(); return !blacklistedIds.contains(userMessage.getId().toString()); } }
注意:如果你的UserMessage的id是数值类型,一定要转成字符串再和黑名单匹配,避免类型不匹配导致判断失效。
方式2:通过Camel Exchange获取缓存(无Spring依赖场景)
如果不想依赖Spring注入,也可以在Bean方法中传入Exchange对象,从Camel上下文直接获取Caffeine组件的缓存:
public boolean isAuthorizedUser(UserMessage userMessage, Exchange exchange) { // 获取Caffeine组件实例 CaffeineCacheComponent caffeineComponent = exchange.getContext().getComponent("caffeine-cache", CaffeineCacheComponent.class); // 获取指定缓存 CaffeineCache cache = caffeineComponent.getCache("blacklistedIds"); List<String> blacklistedIds = (List<String>) cache.get("blacklistedIds"); return !blacklistedIds.contains(userMessage.getId().toString()); }
这种方式更贴合Camel原生组件的使用逻辑,不需要依赖Spring的注入机制。
问题3:是否存在更简洁高效的实现方式?
当然有,这里给你几个优化方向:
1. 优化黑名单加载路由的聚合逻辑
你当前用了completionTimeout(500)等待聚合,其实可以改成completionSize(-1)(表示聚合所有拆分后的行),这样不需要等待超时,文件读取完成后立即存入缓存,效率更高:
from(file(...)) .id(getRouteId()) .split().tokenize("\n") .stopOnException() .aggregate(constant(true), new GroupedBodyAggregationStrategy()) .completionSize(-1) // 替换timeout,完成所有行聚合后执行后续逻辑 .setHeader(CaffeineConstants.ACTION, constant(CaffeineConstants.ACTION_PUT)) .setHeader(CaffeineConstants.KEY, constant("blacklistedIds")) .toF("caffeine-cache://%s", "blacklistedIds")
2. 路由内直接用断言过滤,省去额外Bean
如果过滤逻辑不复杂,可以直接在路由里用simple表达式完成过滤,省去单独的Service类,代码更简洁:
from(kafka(...)) .id(getRouteId()) .unmarshal().json(JsonLibrary.Jackson, UserMessage.class) // 直接通过Caffeine组件表达式获取黑名单并过滤 .filter(simple("${body.id.toString()} not in ${caffeine-cache://blacklistedIds?action=GET&key=blacklistedIds}")) .to(output())
3. 增加黑名单动态刷新机制(可选)
如果黑名单需要定期更新且不想重启应用,可以给加载路由加上定时触发(比如用file组件的delay参数设置轮询间隔),实现文件变化后自动刷新缓存:
// 每5分钟检查一次文件变化并刷新缓存 from("file://your-blacklist-path?delay=300000") .id(getRouteId()) .split().tokenize("\n") .stopOnException() .aggregate(constant(true), new GroupedBodyAggregationStrategy()) .completionSize(-1) .setHeader(CaffeineConstants.ACTION, constant(CaffeineConstants.ACTION_PUT)) .setHeader(CaffeineConstants.KEY, constant("blacklistedIds")) .toF("caffeine-cache://%s", "blacklistedIds")
内容的提问来源于stack exchange,提问作者Abdelghani Roussi

