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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 03:30:45