如何基于Chronicle Queue与Excerpt Tailer实现数据实时重放?
Chronicle Queue 实时数据重放需求实现方案
问题描述
我已经能够读取Chronicle Queue的Excerpt Appender中的所有数据片段,但希望通过Excerpt Tailer读取数据时,模拟数据被持久化时的实时场景。例如:某条数据在下午3点被写入,之后10分钟无新数据,读取时需要等待10分钟再读取下一条片段。目前我只能一次性读取所有可用数据,想知道这个需求是否可行,以及该库的哪部分功能可以用来实现这种实时重放。
当前实现代码
public class Replayer implements Runnable { private static final Logger LOGGER = LogManager.getLogger(Replayer.class); private final ChronicleQueue QUEUE; private final ExcerptTailer TAILER; private final String CYCLE_MODE; private final String CYCLE_STR; private final String REPLAY; private final String NETWORK_INTERFACE; public Replayer(Properties prop) { this(prop, prop.getProperty("CHRONICLE_PATH")); } public Replayer(Properties properties, String path) { String PATH = path; CYCLE_MODE = properties.getProperty("CYCLE_MODE"); CYCLE_STR = properties.getProperty("CYCLE"); REPLAY = properties.getProperty("REPLAY"); NETWORK_INTERFACE = properties.getProperty("NETWORK_INTERFACE"); QUEUE = SingleChronicleQueueBuilder.single(PATH).rollCycle(RollCycles.FAST_DAILY).build(); TAILER = QUEUE.createTailer(); } private MoldUdpHeader moldUdpHeader = new MoldUdpHeader(); private final ByteBuffer byteBuffer = ByteBuffer.allocate(1500); private final byte[] remoteAdd = new byte[4]; private int remotePort; private int remaining; private final ReadBytesMarshallable marshallable = (b) -> { b.read(remoteAdd); remotePort = b.readInt(); remaining = b.readInt(); byteBuffer.put(b.bytesForRead().toByteArray(), 0, remaining); }; private final HashMap<Pair<String, String>, UDPClient> mapping = new HashMap<>(); private final Pair<String, String> pair = new Pair<>("", ""); @Override public void run() { moveToCycle(); System.out.println(currentCycle); //if cycle is not available it prints Integer.MIN_VALUE -2147483648 int counter =0; while (TAILER.readBytes(marshallable)) { if (checkStop()) { break; } pair.setType1(IntStream.range(0, remoteAdd.length).mapToObj(i -> String.valueOf((remoteAdd[i] & 0xFF))).collect(Collectors.joining("."))); pair.setType2(String.valueOf(remotePort)); // System.out.println(IntStream.range(0, remoteAdd.length).mapToObj(i -> String.valueOf((remoteAdd[i] & 0xFF))).collect(Collectors.joining("."))); // System.out.println(" " + remotePort+" "+remaining); if (!mapping.containsKey(pair)) { try { mapping.put(pair.copy(), new UDPClient(remotePort-5000, 1500, remoteAdd, NETWORK_INTERFACE)); } catch (IOException e) { LOGGER.warn(mainMarker, e.getMessage()); } } if (REPLAY.equals("TRUE")) { try { byteBuffer.flip(); mapping.get(pair).sendMessage(byteBuffer); if ((counter = counter % 3) == 0) Thread.sleep(1); //no seq gap } catch (IOException | InterruptedException e) { LOGGER.warn(mainMarker, e.getMessage()); } } moldUdpHeader = (MoldUdpHeader) moldUdpHeader.decode(byteBuffer, 0); System.out.println(moldUdpHeader); byteBuffer.clear(); counter++; } TAILER.close(); QUEUE.close(); } private int currentCycle; private int cycleCounter; public void moveToCycle() { long time; if (CYCLE_STR.equals("TODAY")) { time = TimeUnit.MILLISECONDS.toDays(System.currentTimeMillis()); } else { time = LocalDate.parse(CYCLE_STR, DateTimeFormatter.ofPattern("yyyyMMdd", Locale.US)).toEpochDay(); } if (TAILER.moveToCycle((int) time)) { currentCycle = (int) time; if (!CYCLE_MODE.equals("END")) { try { cycleCounter = currentCycle + Integer.parseInt(CYCLE_MODE); } catch (NumberFormatException e) { LOGGER.warn(mainMarker, e.getMessage()); } } } else { LOGGER.info(mainMarker, "failed to move to specified cycle. Moving to most recent cycle instead"); TAILER.toEnd(); TAILER.moveToCycle(TAILER.cycle()); currentCycle = TAILER.cycle(); } } public boolean checkStop() { int cycle = TAILER.cycle(); if (CYCLE_MODE.equals("END")) { return false; } else { return cycle == cycleCounter; } } }
实现方案
这个需求完全可行,核心思路是记录每条数据的写入时间戳,在重放时计算相邻两条数据的时间差,通过休眠模拟真实的时间间隔。具体实现步骤如下:
1. 写入数据时追加时间戳
在使用Excerpt Appender写入业务数据前,先写入该条数据的创建时间戳(如System.currentTimeMillis()),确保重放时能获取到原始写入时间。示例写入逻辑:
try (ExcerptAppender appender = queue.acquireAppender()) { appender.writeBytes(out -> { out.writeLong(System.currentTimeMillis()); // 写入毫秒级时间戳 // 写入原有业务数据 out.write(remoteAdd); out.writeInt(remotePort); out.writeInt(remaining); out.write(byteBuffer.array(), 0, remaining); }); }
2. 重放时计算时间差并休眠
修改Replayer类的读取逻辑,先解析每条数据的时间戳,再计算与上一条数据的时间差,调用休眠方法模拟真实间隔。核心修改后的代码片段:
private long lastTimestamp = -1; @Override public void run() { moveToCycle(); System.out.println(currentCycle); int counter = 0; // 修改marshallable,先读取时间戳 final ReadBytesMarshallable timeAwareMarshaller = (b) -> { long timestamp = b.readLong(); // 先读取时间戳 // 读取原有业务数据 b.read(remoteAdd); remotePort = b.readInt(); remaining = b.readInt(); byteBuffer.put(b.bytesForRead().toByteArray(), 0, remaining); // 存储时间戳(可以用成员变量传递) this.currentTimestamp = timestamp; }; while (TAILER.readBytes(timeAwareMarshaller)) { if (checkStop()) { break; } // 计算并执行休眠 if (lastTimestamp != -1) { long sleepTime = currentTimestamp - lastTimestamp; if (sleepTime > 0) { try { Thread.sleep(sleepTime); } catch (InterruptedException e) { Thread.currentThread().interrupt(); LOGGER.warn("重放休眠被中断", e); break; } } } lastTimestamp = currentTimestamp; // 原有业务逻辑保持不变 pair.setType1(IntStream.range(0, remoteAdd.length).mapToObj(i -> String.valueOf((remoteAdd[i] & 0xFF))).collect(Collectors.joining("."))); pair.setType2(String.valueOf(remotePort)); if (!mapping.containsKey(pair)) { try { mapping.put(pair.copy(), new UDPClient(remotePort-5000, 1500, remoteAdd, NETWORK_INTERFACE)); } catch (IOException e) { LOGGER.warn(mainMarker, e.getMessage()); } } if (REPLAY.equals("TRUE")) { try { byteBuffer.flip(); mapping.get(pair).sendMessage(byteBuffer); } catch (IOException e) { LOGGER.warn(mainMarker, e.getMessage()); } } moldUdpHeader = (MoldUdpHeader) moldUdpHeader.decode(byteBuffer, 0); System.out.println(moldUdpHeader); byteBuffer.clear(); counter++; } TAILER.close(); QUEUE.close(); }
3. 补充说明
- 若重放的历史数据未包含时间戳,可借助Chronicle Queue的循环文件时间戳估算(比如按日循环的文件,取文件创建时间),但精度会受限。
- 休眠时需处理
InterruptedException,保证线程中断时能优雅退出。 - 如需更高精度的时间控制,可使用
LockSupport.parkNanos()替代Thread.sleep()。
内容的提问来源于stack exchange,提问作者chunkynuggy
相关产品推荐
相关产品推荐

