Quarkus缓存Pub/Sub订阅数据30分钟的实现与验证方法咨询
实现方法
1. 添加必要依赖
首先引入Quarkus缓存扩展,若需要分布式缓存(多实例场景),可搭配Redis缓存后端;同时引入Pub/Sub订阅扩展:
<!-- 基础缓存依赖 --> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-cache</artifactId> </dependency> <!-- 若用Redis作为缓存后端,添加此依赖 --> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-redis-cache</artifactId> </dependency> <!-- Pub/Sub订阅依赖 --> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-google-cloud-pubsub</artifactId> </dependency>
2. 配置缓存过期时间
在application.properties中设置缓存过期规则(30分钟=1800秒),同时配置Pub/Sub订阅信息:
# Caffeine本地缓存默认过期时间 quarkus.cache.caffeine.default.expire-after-write=1800s # Redis缓存配置(仅使用Redis时需要) quarkus.redis.hosts=redis://localhost:6379 quarkus.cache.redis.default.expire-after-write=1800s # Pub/Sub订阅配置 mp.messaging.incoming.pubsub-subscription.subscription-id=你的订阅ID mp.messaging.incoming.pubsub-subscription.project-id=你的GCP项目ID
3. 实现Pub/Sub消息消费与缓存逻辑
创建订阅消费者类,接收消息后将数据存入缓存。注意用消息的唯一标识作为缓存Key(避免数据覆盖):
import io.quarkus.cache.CacheResult; import io.smallrye.reactive.messaging.annotations.Blocking; import org.eclipse.microprofile.reactive.messaging.Incoming; import jakarta.enterprise.context.ApplicationScoped; @ApplicationScoped public class PubSubCacheConsumer { private static final String PUBSUB_CACHE = "pubsub-data-cache"; @Incoming("pubsub-subscription") @Blocking // 缓存操作可能阻塞,根据实际场景调整 public void consumeMessage(String payload) { String cacheKey = extractUniqueKey(payload); cacheData(cacheKey, payload); } // @CacheResult自动将返回值存入指定缓存,过期时间遵循配置 @CacheResult(cacheName = PUBSUB_CACHE) public String cacheData(String key, String data) { return data; } // 从消息体提取唯一Key(示例:若为JSON则解析ID字段,需根据实际消息格式调整) private String extractUniqueKey(String payload) { return payload; // 简化处理,实际需替换为真实唯一标识 } }
4. 缓存读取逻辑(可选)
若需要读取缓存数据,添加如下方法:
@CacheResult(cacheName = PUBSUB_CACHE) public String getCachedData(String key) { // 缓存未命中时返回null,可根据需求添加兜底逻辑 return null; }
验证缓存时长的方法
1. 单元测试验证
编写测试类,模拟消息存入缓存,等待指定时间后检查数据是否过期:
import io.quarkus.test.junit.QuarkusTest; import jakarta.inject.Inject; import org.junit.jupiter.api.Test; import java.util.concurrent.TimeUnit; import static org.junit.jupiter.api.Assertions.*; @QuarkusTest public class CacheExpiryTest { @Inject PubSubCacheConsumer consumer; @Test public void testCacheExpiry() throws InterruptedException { String testKey = "test-unique-key"; String testData = "test-payload"; // 存入缓存 consumer.cacheData(testKey, testData); // 立即校验缓存存在 assertEquals(testData, consumer.getCachedData(testKey)); // 等待30分钟10秒,校验缓存过期 TimeUnit.MINUTES.sleep(30); TimeUnit.SECONDS.sleep(10); assertNull(consumer.getCachedData(testKey)); } }
2. 工具命令验证
- Redis缓存:直接使用Redis CLI查看Key的剩余过期时间:
redis-cli TTL pubsub-data-cache::test-unique-key
返回值初始接近1800秒,随时间递减,过期后返回-2。
- Caffeine本地缓存:启用Quarkus Metrics扩展(
quarkus-micrometer),通过Prometheus查询cache_expirations_total指标,确认缓存过期记录。
3. 日志验证
在缓存方法中添加日志,记录存入/读取时间与缓存状态:
import org.jboss.logging.Logger; private static final Logger LOG = Logger.getLogger(PubSubCacheConsumer.class); @CacheResult(cacheName = PUBSUB_CACHE) public String cacheData(String key, String data) { LOG.info(String.format("存入缓存Key: %s,时间戳: %d", key, System.currentTimeMillis())); return data; } @CacheResult(cacheName = PUBSUB_CACHE) public String getCachedData(String key) { LOG.info(String.format("读取缓存Key: %s,时间戳: %d", key, System.currentTimeMillis())); return null; }
观察日志:存入30分钟后读取时,若缓存过期会重新执行cacheData方法并打印存入日志,或读取返回null。
内容的提问来源于stack exchange,提问作者SoftCore Productions
相关产品推荐
相关产品推荐

