如何配置StoreFileListener实现Chronicle Queue的数据留存
Chronicle Queue 中 StoreFileListener 的使用与定时删除文件实现
Chronicle Queue 官方文档提到可以通过 StoreFileListener 监听队列文件的添加或释放事件,以此实现过期数据文件的定时删除。下面是具体的实现思路和示例代码:
核心思路
StoreFileListener 提供了两个关键回调方法:
onReleased(File file):当某个队列文件不再被读写使用时触发onAcquired(File file):当队列文件被创建或重新启用时触发
我们可以在 onReleased 中记录下已释放的文件,再通过定时任务定期检查这些文件是否超过设定的保留时长,若超过则执行删除操作。
完整示例代码
import net.openhft.chronicle.queue.ChronicleQueue; import net.openhft.chronicle.queue.StoreFileListener; import net.openhft.chronicle.queue.impl.single.SingleChronicleQueueBuilder; import java.io.File; import java.util.ArrayList; import java.util.List; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; public class QueueFileCleanerExample { // 数据文件保留时长(示例:24小时) private static final long RETENTION_DURATION = 24 * 60 * 60 * 1000; // 定时清理任务的执行间隔(示例:每小时执行一次) private static final long CLEANUP_INTERVAL = 60 * 60 * 1000; // 存储已释放的待清理文件 private final List<File> releasedFiles = new ArrayList<>(); private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); public static void main(String[] args) { QueueFileCleanerExample example = new QueueFileCleanerExample(); example.startQueueAndCleaner(); } private void startQueueAndCleaner() { // 创建Chronicle Queue实例 try (ChronicleQueue queue = SingleChronicleQueueBuilder.binary(new File("./queue-data")) // 注册自定义的StoreFileListener .storeFileListener(new CustomStoreFileListener()) .build()) { // 启动定时清理任务 scheduler.scheduleAtFixedRate(this::cleanupExpiredFiles, 0, CLEANUP_INTERVAL, TimeUnit.MILLISECONDS); // 这里可以添加队列的读写逻辑,示例省略 // ... // 保持程序运行,实际场景中根据业务需求控制生命周期 Thread.currentThread().join(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { scheduler.shutdown(); } } // 自定义StoreFileListener实现 private class CustomStoreFileListener implements StoreFileListener { @Override public void onReleased(File file) { // 文件释放时加入待清理列表 synchronized (releasedFiles) { releasedFiles.add(file); } System.out.println("文件已释放,加入待清理列表: " + file.getAbsolutePath()); } @Override public void onAcquired(File file) { // 如果已释放的文件被重新启用,从待清理列表移除 synchronized (releasedFiles) { releasedFiles.remove(file); } System.out.println("文件被重新启用,从待清理列表移除: " + file.getAbsolutePath()); } } // 清理过期文件的方法 private void cleanupExpiredFiles() { long currentTime = System.currentTimeMillis(); synchronized (releasedFiles) { // 遍历待清理列表,删除超过保留时长的文件 releasedFiles.removeIf(file -> { long lastModified = file.lastModified(); if (currentTime - lastModified > RETENTION_DURATION) { boolean deleted = file.delete(); if (deleted) { System.out.println("已删除过期文件: " + file.getAbsolutePath()); } else { System.err.println("删除文件失败: " + file.getAbsolutePath()); } return deleted; } return false; }); } } }
关键说明
- 线程安全:由于
StoreFileListener的回调和定时清理任务可能在不同线程执行,所以对releasedFiles的操作需要加锁同步 - 保留时长与清理间隔:可以根据业务需求调整
RETENTION_DURATION和CLEANUP_INTERVAL的值 - 文件重新启用处理:如果已释放的文件被重新读取(比如回溯消费),
onAcquired会触发,此时需要从待清理列表移除该文件,避免误删 - 资源关闭:程序退出时要关闭
ScheduledExecutorService和ChronicleQueue,避免资源泄漏
内容的提问来源于stack exchange,提问作者Jimmy Guo
相关产品推荐
相关产品推荐

