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

Aeron Archive同一录制场景下recordingId为何自动递增?

关于Aeron Archive recordingId递增及多录制文件的问题

基于aeroncookbook.com示例搭建ArchiveHost,该主机包含一个被录制的Publication,同时将Subscription接收的所有消息转发至这个被录制的Publication,核心代码如下:

private void createArchiveAndRecord() {

    final var archiveContext = new Archive.Context()
            .deleteArchiveOnStart(true)
            .replicationChannel(REPLICATION_CHANNEL)
            .controlChannel(ARCHIVE_CONTROL_REQUEST_CHANNEL)
            .controlStreamId(ARCHIVE_CONTROL_REQUEST_STREAM_ID)
            .threadingMode(ArchiveThreadingMode.DEDICATED)
            .idleStrategySupplier(BackoffIdleStrategy::new)
            .errorHandler(this::errorHandler);

    var mediaDriverCtx = new MediaDriver.Context()
            .errorHandler(this::errorHandler)
            .threadingMode(ThreadingMode.DEDICATED)
            .conductorIdleStrategy(new BackoffIdleStrategy())
            .receiverIdleStrategy(new BackoffIdleStrategy())
            .senderIdleStrategy(new BackoffIdleStrategy())
            .dirDeleteOnStart(true);

    archivingMediaDriver = ArchivingMediaDriver.launch(
            mediaDriverCtx,
            archiveContext
    ); 

    try {

         aeron = Aeron.connect(
                 new Aeron.Context()
                         .aeronDirectoryName(archivingMediaDriver.mediaDriver().aeronDirectoryName())
                         .errorHandler(this::errorHandler)
         );
         archive = AeronArchive.connect(
                 new AeronArchive.Context()
                         .aeron(aeron)
                         .controlRequestChannel(ARCHIVE_CONTROL_REQUEST_CHANNEL)
                         .controlResponseChannel(ARCHIVE_CONTROL_RESPONSE_CHANNEL)
                         .controlRequestStreamId(ARCHIVE_CONTROL_REQUEST_STREAM_ID)
                         .recordingSignalConsumer(new ArchiveActivityListener())
                         .errorHandler(this::errorHandler)
         );

        archive.startRecording(RECORDING_CHANNEL, STREAM_ID, SourceLocation.REMOTE);
        final Publication pub = aeron.addPublication(PUBLICATION_CHANNEL, STREAM_ID);
        final CountersReader counters = aeron.countersReader();
        final int recordingCounterId = awaitRecordingCounterId(counters, pub.sessionId());
        final long recordingId = RecordingPos.getRecordingId(counters, recordingCounterId);

        this.currentState = State.RDY;
        this.subscriptionFragmentHandler.setPublication(pub);
    } catch (InterruptedException e) {
        throw new RuntimeException(e);
    }
}

有5个行情数据生产者向该ArchiveHost的Subscription发送消息,整体运行正常,但recordingId会随机递增,生成了多个录制文件:

-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 18:54 0-0.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:20 0-1073741824.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:22 0-1207959552.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 18:58 0-134217728.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:23 0-1342177280.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:24 0-1476395008.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:26 0-1610612736.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:28 0-1744830464.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:30 0-1879048192.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:33 0-2013265920.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:36 0-2147483648.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:38 0-2281701376.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:38 0-2415919104.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:01 0-268435456.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:04 0-402653184.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:08 0-536870912.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:12 0-671088640.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:15 0-805306368.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:18 0-939524096.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:40 1-2432696320.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:41 1-2566914048.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:43 1-2701131776.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:43 2-2785017856.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:45 3-2801795072.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:48 3-2936012800.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:50 3-3070230528.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:51 3-3204448256.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:53 4-3221225472.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 19:56 4-3355443200.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:00 4-3489660928.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:02 4-3623878656.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:04 4-3758096384.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:06 5-3858759680.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:07 5-3992977408.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:09 5-4127195136.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:10 5-4261412864.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 20:11 5-4395630592.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 22:21 6-10032775168.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 22:23 6-10166992896.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 22:26 6-10301210624.rec
-rw-r--r--. 1 ec2-user ec2-user 134217728 Jun 11 22:30 6-10435428352.rec

ArchiveHost从未停止,原本预期运行期间recordingId保持不变,请教:

  1. 这种情况是否正常?
  2. 可能的原因是什么?
  3. 能否添加监听或日志来获取recordingId递增时的详细信息?

问题解答

1. 这种情况是否正常?

这种情况不属于预期的正常行为。正常情况下,针对同一个CHANNEL和STREAM_ID的录制,只要Archive持续运行且录制未被手动停止/重启,应该保持同一个recordingId,生成连续的录制文件(文件名中后缀的递增是文件分段,而非recordingId变更)。你看到的文件名前半段数字(如0、1、2...)是recordingId,后半段是文件起始位置,当前recordingId频繁变更属于异常。

2. 可能的原因是什么?

  • 录制会话意外中断并自动重启:Aeron Archive会在录制会话因某些异常中断后自动尝试重启录制,比如:
    • 转发的Publication出现异常(如会话超时、连接中断),导致录制的数据源丢失,Archive自动重启录制;
    • 5个生产者的消息发送出现异常波动,导致Archive的录制订阅出现短暂的无消息超时,触发录制重启;
    • Archive内部的录制计数器或元数据出现异常,导致系统认为当前录制已终止,启动新的录制。
  • 重复调用startRecording:检查代码中是否存在重复调用archive.startRecording的逻辑,比如在消息转发的处理逻辑中意外触发了多次录制启动请求,每次调用都会生成新的recordingId。
  • 配置参数问题:如果Archive的配置中设置了自动分段的阈值(如maxRecordingLength),但该阈值设置过小,会导致录制频繁分段,但分段不会变更recordingId,不过如果分段过程中出现异常,可能触发新的录制会话。
  • 媒体驱动或Archive的线程问题:线程池资源不足、IdleStrategy配置不合理导致线程无法及时处理事件,引发录制会话的异常中断。

3. 添加监听或日志获取详细信息的方法

  • 完善ArchiveActivityListener:你已经设置了recordingSignalConsumer(new ArchiveActivityListener()),可以扩展这个监听器的逻辑,监听RecordingSignal中的事件,比如START、STOP、EXTEND等,记录每个事件对应的recordingId、时间戳和事件原因:
class ArchiveActivityListener implements RecordingSignalConsumer {
    @Override
    public void onSignal(long recordingId, long subscriptionId, long sessionId, RecordingSignal signal, String channel, int streamId, long position) {
        // 打印详细日志
        System.out.printf("Recording Signal: signal=%s, recordingId=%d, streamId=%d, position=%d, channel=%s%n",
                signal.name(), recordingId, streamId, position, channel);
        // 可以写入日志文件或监控系统
    }
}
  • 增强ErrorHandler日志:在自定义的errorHandler中,详细记录所有Archive和MediaDriver的错误信息,包括异常类型、堆栈和上下文,很多录制异常会通过errorHandler抛出:
private void errorHandler(Throwable throwable) {
    // 使用日志框架(如SLF4J)记录详细错误
    LoggerFactory.getLogger(getClass()).error("Aeron/Archive Error occurred", throwable);
}
  • 监控录制计数器:通过CountersReader定期检查当前活跃的录制计数器,记录recordingId的变化,比如定时轮询:
// 定时任务中执行
CountersReader counters = aeron.countersReader();
RecordingPos.forEachRecording(counters, (counterId, recordingIdValue, position) -> {
    System.out.printf("Active Recording: recordingId=%d, position=%d%n", recordingIdValue, position);
});
  • 启用Archive的内置日志:在Archive.Context中配置日志级别,开启更详细的Archive内部日志,比如:
archiveContext.logger(LoggerFactory.getLogger("io.aeron.archive"));

确保日志框架配置为DEBUG级别,这样可以捕获Archive内部的录制启动、停止等关键操作日志。


内容的提问来源于stack exchange,提问作者Lauda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 11:19:54