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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 02:57:16