news 2026/9/25 8:03:40

Aeron Samples 实战指南:从 Basic Pub/Sub 到吞吐量与延迟测试的完整工具箱

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Aeron Samples 实战指南:从 Basic Pub/Sub 到吞吐量与延迟测试的完整工具箱
  • 消息队列
  • 后端
  • 通信

【免费下载链接】aeron

Efficient reliable UDP unicast, UDP multicast, and IPC message transport

项目地址:https://gitcode.com/gh_mirrors/ae/aeron
点击查看免费下载

本指南以 aeron-samples/README.md 为骨架,系统讲解 Aeron 官方示例与工具集的用法:先介绍如何启动媒体驱动(Media Driver),再逐一演示 BasicSubscriber、BasicPublisher、RateSubscriber、StreamingPublisher 等消息收发样例,接着深入 AeronStat、ErrorStat、LossStat、BacklogStat、LogInspector 等监控诊断工具,最后覆盖 Archive 录播样例、嵌入式吞吐/延迟测试以及跨进程网络测试。读完本文,你将掌握在真实环境中搭建 Aeron 收发链路、观测系统状态并做性能基准测试的完整实战能力。

一、运行前置条件:先启动媒体驱动(Media Driver)

Aeron 采用"客户端 + 媒体驱动"的进程模型:客户端应用通过共享内存中的命令与控制(CnC)文件与媒体驱动通信,实际的数据收发、重传、流量控制等工作全部由媒体驱动完成。因此,运行任何样例之前,必须先在独立的控制台中启动媒体驱动,命令如下(来自 aeron-samples/README.md):

aeron-samples/scripts/media-driver <optional properties file> aeron-samples/scripts/low-latency-media-driver <optional properties file>
  • media-driver脚本本质上等价于执行 run-java 启动io.aeron.driver.MediaDriver,即标准配置的媒体驱动;
  • low-latency-media-driver会额外加载low-latency.properties配置(见 low-latency-media-driver),适合对时延敏感的场景。

<optional properties file>可以是一个 Java properties 文件,用于覆盖媒体驱动的各项参数(如线程模式、MTU、socket 缓冲区等),不传则使用默认配置。所有运行脚本都依赖 java-common 中定义的环境:

  • 必须设置JAVA_HOME环境变量,否则脚本会报错退出;
  • 脚本从仓库根目录的 version.txt 读取版本号,然后加载aeron-all/build/libs/aeron-all-${VERSION}.jar作为 classpath;
  • 默认附带了面向高性能的 JVM 参数:-XX:+TrustFinalNonStaticFields、-XX:GuaranteedSafepointInterval=300000(降低安全点频率)、-XX:+UseParallelGC,以及为最新 JDK 准备的--add-opens模块开关;
  • 额外 JVM 参数可通过JVM_OPTS环境变量注入,例如JVM_OPTS="-Daeron.sample.channel=..." ./script。

嵌入式媒体驱动模式

从源码看,多数样例类(如 BasicPublisher.java 与 BasicSubscriber.java)都读取SampleConfiguration.EMBEDDED_MEDIA_DRIVER(对应-Daeron.sample.embeddedMediaDriver=true系统属性)。设置为true时,应用会通过MediaDriver.launchEmbedded()在同一进程内启动媒体驱动,并把aeronDirectoryName指向该驱动,从而无需先启动独立媒体驱动进程。默认值为false,即默认依赖外部媒体驱动。

二、SampleConfiguration:所有样例共享的配置中心

所有样例的默认参数都定义在 SampleConfiguration.java 中,并且全部支持通过 Java 系统属性覆盖。核心参数如下表:

系统属性默认值说明
aeron.sample.channelaeron:udp?endpoint=localhost:20121默认样例使用的 UDP 通道 URI
aeron.sample.streamId1001默认流 ID
aeron.sample.ping.channelaeron:udp?endpoint=localhost:20123Ping 样例发送通道
aeron.sample.pong.channelaeron:udp?endpoint=localhost:20124Pong 样例响应通道
aeron.sample.ping.streamId1002Ping 流 ID
aeron.sample.pong.streamId1003Pong 流 ID
aeron.sample.messages10000000发送/测试的消息条数
aeron.sample.messageLength32消息长度(字节)
aeron.sample.warmup.messages10000预热消息数
aeron.sample.warmup.iterations10预热迭代次数
aeron.sample.frameCountLimit10每次 poll 最多处理的消息数
aeron.sample.lingerTimeout0发送完成后的驻留时间(毫秒)
aeron.sample.idleStrategyorg.agrona.concurrent.BusySpinIdleStrategy空闲策略类名
aeron.sample.embeddedMediaDriverfalse是否内嵌媒体驱动
aeron.sample.exclusive.publicationsfalse是否使用独占 Publication
aeron.sample.randomMessageLengthfalse是否随机消息长度

例如,要修改 BasicSubscriber 的通道与流 ID,可以这样启动:

aeron-samples/scripts/basic-subscriber -Daeron.sample.channel=aeron:udp?endpoint=localhost:5555 -Daeron.sample.streamId=20

注意:JAVA_OPTS不会直接生效,脚本只透传JVM_OPTS与命令行位置参数。若通过脚本传系统属性,请使用JVM_OPTS环境变量,例如:

JVM_OPTS="-Daeron.sample.channel=aeron:udp?endpoint=localhost:5555 -Daeron.sample.streamId=20" ./aeron-samples/scripts/basic-subscriber

三、入门收发样例:BasicPublisher 与 BasicSubscriber

这是官方文档列出的第一组样例,用于快速验证 Aeron 端到端消息通路。

BasicPublisher

BasicPublisher.java 的职责:向指定通道和流 ID 发布固定数量的消息,每条消息之间暂停 1 秒,发送完成后驻留LINGER_TIMEOUT_MS毫秒,给可能发生丢失的订阅方留出通过 NAK 请求重传并恢复数据的时间。核心代码逻辑为:

  1. 通过Aeron.connect(ctx)连接媒体驱动,aeron.addPublication(CHANNEL, STREAM_ID)创建发布端;
  2. 将"Hello World! " + i写入UnsafeBuffer,调用publication.offer(buffer, 0, length)发送;
  3. 根据返回值判断发送结果:
    • 返回正值:发送成功(已推进到的位置);
    • Publication.BACK_PRESSURED:背压导致 offer 失败;
    • Publication.NOT_CONNECTED:尚无订阅方连接;
    • Publication.ADMIN_ACTION:系统管理操作期间失败;
    • Publication.CLOSED:Publication 已关闭;
    • Publication.MAX_POSITION_EXCEEDED:达到最大位置。
  4. 每次 offer 后调用publication.isConnected()检测是否有活动订阅方。

运行方式:

aeron-samples/scripts/basic-publisher

对应脚本 basic-publisher 实际执行io.aeron.samples.BasicPublisher主类。

BasicSubscriber

BasicSubscriber.java 的职责:订阅指定通道与流 ID,将收到的每条消息(或消息片段)以 ASCII 形式打印出来。它注册了availableImageHandler与unavailableImageHandler回调来打印映像(Image)的可用/不可用事件,然后通过SamplesUtil.subscriberLoop(fragmentHandler, FRAGMENT_COUNT_LIMIT, running)进入轮询循环,每次 poll 最多处理FRAGMENT_COUNT_LIMIT(默认 10)条消息。

运行方式:

aeron-samples/scripts/basic-subscriber

对应脚本 basic-subscriber 实际执行io.aeron.samples.BasicSubscriber主类。该应用只处理非分片消息;如需处理大消息分片重组,README 与类注释都指向带FragmentAssembler的变体实现。

体验步骤

  1. 终端 A:aeron-samples/scripts/media-driver;
  2. 终端 B:aeron-samples/scripts/basic-subscriber;
  3. 终端 C:aeron-samples/scripts/basic-publisher;
  4. 观察终端 B 打印出 "Hello World! 0"、"Hello World! 1"……,终端 C 打印 "Offering i/N - yay!"。

四、速率观测样例:RateSubscriber 与 StreamingPublisher

  • RateSubscriber:一个按速率打印消息接收情况的订阅方。它借助 ImageRateReporter 统计单位时间内接收的消息数与吞吐,并通过ImageRateSubscriber主类运行,适合观察持续流的接收速率。
  • StreamingPublisher:以"尽可能快"的方式流式发布消息,并实时展示发布速率,与 RateSubscriber 搭配可以直观地看到 Aeron 在 UDP 通道上的线速收发能力。

两者对应的脚本为 rate-subscriber 与 streaming-publisher。此外仓库还提供streaming-exclusive-publisher(C 语言实现)与streaming-exclusive-publisher相关 Java 变体,以及 C 语言的 rate_subscriber.c 与 streaming_publisher.c,方便对照跨语言行为。

五、Ping/Pong:延迟测试工具

Ping/Pong 是经典的往返延迟(RTT)测量模型:

  • Ping:Ping 侧发送端。脚本 ping 执行io.aeron.samples.Ping,默认参数为 100 万条消息、消息长度 32 字节,并开启-Dagrona.disable.bounds.checks=true(关闭边界检查)、-Daeron.pre.touch.mapped.memory=true(预触达映射内存)、-Daeron.sample.exclusive.publications=true(使用独占发布端);
  • Pong:Pong 侧响应端。脚本 pong 执行io.aeron.samples.Pong,将 Ping 发来的消息原样回送,同样开启上述高性能选项。

二者使用独立的通道(localhost:20123/localhost:20124)与流 ID(1002 / 1003),在 Ping 侧用 HdrHistogram 统计 RTT 分布。Java 实现位于io.aeron.samples.Ping/io.aeron.samples.Pong,C 语言对应物为 cping.c 与 cpong.c。运行方法:先启动媒体驱动,再分别启动pong与ping,观察 Ping 侧输出的延迟直方图。

六、性能基准测试:嵌入式与跨进程变体

官方文档说明,仓库还附带了一组性能测试,既可以在同一进程内运行(方便起见,无需媒体驱动),也可以跨进程运行,并且每种都提供了吞吐量与延迟测量的变体:

  • embedded:测试倾向于在同一进程内运行。例如 EmbeddedThroughput.java(吞吐量)、EmbeddedPingPong.java(延迟)、EmbeddedIpcThroughput.java、EmbeddedExclusiveThroughput.java、EmbeddedExclusiveSpiedThroughput.java、EmbeddedBufferClaimIpcThroughput.java、EmbeddedExclusiveVectoredIpcThroughput.java 等;
  • media:面向 IPC(共享内存)或 UDP(网络)的变体,需要在独立媒体驱动下跨进程运行。

以 embedded-throughput 脚本为例,它执行io.aeron.samples.EmbeddedThroughput并预设了一组对吞吐友好的参数,可以直接照抄到自己的基准测试中:

exec "${DIR}/run-java" \ -Djava.net.preferIPv4Stack=true \ -Dagrona.disable.bounds.checks=true \ -Daeron.sample.messageLength=32 \ -Daeron.sample.messages=500000000 \ -Daeron.term.buffer.sparse.file=false \ -Daeron.mtu.length=8k \ -Daeron.socket.so_sndbuf=2m \ -Daeron.socket.so_rcvbuf=2m \ -Daeron.rcv.initial.window.length=2m \ ${JVM_OPTS} io.aeron.samples.EmbeddedThroughput "$@"

这些参数的含义:aeron.mtu.length=8k提高 MTU 以减少报文数量,aeron.socket.so_sndbuf/so_rcvbuf=2m加大 socket 缓冲区,aeron.rcv.initial.window.length=2m扩大接收窗口以容忍更多在途数据,aeron.term.buffer.sparse.file=false关闭稀疏 term 文件。

从 EmbeddedPingPong.java 的源码可以看到延迟测试的实现要点:它使用ThreadingMode.DEDICATED专用线程模式,conductor 用BackoffIdleStrategy,sender/receiver 用NoOpIdleStrategy(纯忙等,牺牲 CPU 换取最低延迟),消息用UnsafeBuffer直接分配并缓存行对齐(BitUtil.CACHE_LINE_LENGTH),结果存入 HdrHistogram 的Histogram,同时支持WARMUP_NUMBER_OF_ITERATIONS次预热迭代。C++ 侧对应实现见 PingPong.cpp、ExclusivePingPong.cpp、Throughput.cpp、ExclusiveThroughput.cpp,C 语言侧有 ping_pong_raw.c 等。

七、监控与诊断工具:AeronStat / ErrorStat / LossStat / BacklogStat / LogInspector

官方文档明确列出以下监控与诊断工具:

工具脚本用途
AeronStataeron-stat打印媒体驱动正在使用的计数器(counter)的标签与数值
ErrorStaterror-stat打印媒体驱动观测到的去重后的错误(distinct errors)
LossStatloss-stat按流打印丢包报告
BacklogStatbacklog-stat打印各流的流位置报告,给出每个流的处理积压(backlog)指示
LogInspectorlog-inspector诊断工具,打印指定流的日志缓冲区(log buffer)内容,用于调试

AeronStat 详解

AeronStat.java 的原理:媒体驱动会在共享内存中维护一个命令与控制(CnC)文件(布局见CncFileDescriptor),AeronStat 读取该文件并通过CountersReader打印所有计数器,默认每秒刷新一次并持续监听(watch模式),也支持命令行过滤参数:

java -cp aeron-samples/build/libs/samples.jar io.aeron.samples.AeronStat type=[1-9] identity=12345

支持的过滤键包括type(计数器类型)、identity(系统计数器 ID 或位置计数器的注册 ID)、session、stream、channel。其中计数器类型划分:0 为系统计数器,1–5、9、10、11 为流位置与指示器(如PUBLISHER_POS、SENDER_LIMIT、RECEIVER_POS、PER_IMAGE等),6–7 为通道端点状态(SEND_CHANNEL_STATUS、RECEIVE_CHANNEL_STATUS)。这些类型常量定义在io.aeron.driver.status包中(如PublisherPos.PUBLISHER_POS_TYPE_ID、PublisherLimit.PUBLISHER_LIMIT_TYPE_ID),运行时可通过delay与watch参数控制刷新间隔与持续监听行为。

LogInspector 与 LogInspectorCli

LogInspector 用于离线检查某个流日志缓冲区的完整内容(帧头、位置、数据),对排查发布端/接收端位置不一致、帧损坏等问题很有帮助;仓库还提供命令行交互版本 LogInspectorCli.java 与对应的 log-inspector-cli 脚本。C 语言侧对应工具为 driver_tool.c 与 aeron_stat.c,可用 CMake 构建后运行。

八、Aeron Archive 录播样例

官方文档指出,在 aeron-samples/scripts/archive/ 子目录(及 Java 包)中可以找到基于 Archive 的流录制与回放样例,包括录播吞吐测试、录制发布/回放订阅、ReplayMerge 与持久订阅(Persistent Subscription)等。目录中同时提供了.cmd(Windows)与无后缀(Unix)两套脚本,以及三个 properties 模板:standard-archive.properties、high-throughput-archive.properties、lightweight-archive.properties。

录播吞吐测试(Embedded Throughput Samples)

  • embedded-recording-throughput:录制一批消息,然后询问是否重复测试,多次重复可让系统预热。文档建议调整写入同步级别以权衡持久化与性能:
    • aeron.archive.file.sync.level=0:普通写入 OS 页缓存,由后台刷盘;
    • aeron.archive.file.sync.level=1:强制脏数据页落盘;
    • aeron.archive.file.sync.level=2:强制脏数据页与文件元数据落盘。 级别越高,崩溃恢复时数据越安全,但吞吐/时延开销越大。
  • embedded-replay-throughput:先录制一批消息,再在新的流上回放,同样支持重复回放以预热。文档建议尝试不同的消息长度与线程配置。

这两个样例默认会把 Archive 目录建在临时文件系统上,因此强烈建议通过 properties 文件显式指定aeron.archive.dir指向快速存储(如 NVMe/内存盘)。参考 standard-archive.properties 的内容:

aeron.archive.dir=../../build/archive aeron.archive.threading.mode=SHARED aeron.archive.file.sync.level=0 aeron.archive.file.io.max.length=1m aeron.spies.simulate.connection=true aeron.threading.mode=SHARED aeron.term.buffer.sparse.file=true

录制发布 + 回放订阅的完整流程

来自 archive/README.md 的标准操作步骤:

  1. 启动归档媒体驱动(自己的控制台):
    ./archiving-media-driver <config properties file>
  2. 启动带录制功能的发布端:
    ./recorded-basic-publisher <config properties file>
  3. 启动普通订阅方,使发布端连接上并开始录制:
    cd .. && ./basic-subscriber <config properties file>
  4. 启动请求回放的订阅方,消费已录制流:
    ./replay-basic-subscriber <config properties file>
  5. (可选)运行 AeronStat 观察状态,重点关注录制流对应的rec-pos(录制位置)计数器:
    ./aeron-stat
  6. 检查错误:
    ./error-stat

注意:archiving-media-driver、recorded-basic-publisher、replay-basic-subscriber等脚本位于aeron-samples/scripts/archive/目录,而basic-subscriber、aeron-stat、error-stat位于上一级目录aeron-samples/scripts/,执行时需按上述相对路径或使用绝对路径。

ReplayMerge:直播 + 回放的合并

ReplayMerge 适用于 MDC(Multi-Destination Cast)或组播发布场景,先将发布流录制到 Archive,同时让订阅方在直播流与回放流之间无缝合并,适合"边录边看、断点续播"的实时体验。典型步骤:

  1. 终端一:./archiving-media-driver;
  2. 终端二:用 MDC 动态控制通道发布并录制:
    export JVM_OPTS="-Daeron.sample.channel=aeron:udp?control=localhost:20550|control-mode=dynamic|alias=replay-merge-sample" ./recorded-basic-publisher
  3. 终端三:启动与发布通道匹配的普通订阅方:
    export JVM_OPTS="-Daeron.sample.channel=aeron:udp?control=localhost:20550" ../basic-subscriber
  4. 稍等片刻,再启动 ReplayMerge 订阅方,aeron.sample.channel与发布通道保持一致:
    export JVM_OPTS="-Daeron.sample.channel=aeron:udp?control=localhost:20550" ./replay-merge-subscriber

持久订阅(Persistent Subscription)

持久订阅允许订阅方在中间离线一段时间后,从 Archive 中恢复错过的消息,实现"不会丢消息的订阅"。步骤(同样可用JVM_OPTS传系统属性定制配置):

  1. ./archiving-media-driver;
  2. ./recorded-basic-publisher;
  3. ../basic-subscriber(确保录制进行);
  4. 让消息发送一段时间后,启动./persistent-subscriber,它会在回放位置之后从录制流继续消费。

其他 Archive 样例

目录下还有recording-replicator(录制副本复制,用于跨节点/跨区域复制录制数据)、segment-inspector(检查 Archive 分段文件内容)等脚本,以及 archive 样例源码 与aeron-system-tests/src/test/java/io/aeron/archive下的系统测试可供深入学习。

九、其他实用脚本与扩展阅读

  • raw 目录:send-receive-udp-ping/send-receive-udp-pong系列脚本演示在未使用 Aeron 的情况下直接用 UDP 收发报文(配合 ping_pong_raw.c),用于理解底层 socket 行为与对比基线;
  • cluster 目录:提供basic-auction-cluster、basic-auction-client等集群样例启动脚本与 namespace 管理脚本,配套教程见 Cluster-Tutorial.asciidoc;
  • response 目录:response_client.c/response_server.c与 Java 的ResponseClient/ResponseServer演示了 Aeron 的请求-响应(Response Channels)模式;
  • echo / stress / security 包:分别提供 echo 回显、压力测试与安全样例;
  • FileSender / FileReceiver:演示通过 Aeron IPC/UDP 传输文件的场景;
  • C++ 样例:BasicPublisher.cpp、BasicSubscriber.cpp、Ping.cpp、Pong.cpp、RateSubscriber.cpp、StreamingPublisher.cpp等位于 aeron-samples/src/main/cpp,可用 CMake 构建(见 CMakeLists.txt 与 aeron-samples/src/main/c/CMakeLists.txt),与 Java 样例行为一一对应。

十、小结

Aeron Samples 是一个覆盖面完整的官方演练场:从media-driver启动、BasicPublisher/BasicSubscriber入门,到Ping/Pong延迟测试、embedded-*吞吐基准,再到AeronStat/ErrorStat/LossStat/BacklogStat/LogInspector监控诊断,以及 Archive 录制/回放/持久订阅/ReplayMerge 的高级场景,均可在 aeron-samples 模块内开箱即用地复现。所有样例共享SampleConfiguration的默认参数与系统属性覆盖机制,配合JVM_OPTS注入参数,你可以快速在自己的环境下调整通道、流 ID、消息量与消息长度,把官方样例改造成属于自己的压测与观测工具。

  • 消息队列
  • 后端
  • 通信

【免费下载链接】aeron

Efficient reliable UDP unicast, UDP multicast, and IPC message transport

项目地址:https://gitcode.com/gh_mirrors/ae/aeron
点击查看免费下载

相关推荐

上一篇:3分钟掌握Mermaid:用代码绘制专业图表的高效工具
下一篇:Android应用隐私合规检查:UltimateAndroidReference中的完整工具指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 8:03:07

DeepSpeed 高级安装指南:Ops 预编译、多节点分发与 GPU 架构定制

推理引擎大模型 【免费下载链接】FlexGen Running large language models on a single GPU for throughput-oriented scenarios. 项目地址&#xff1a; https://gitcode.com/gh_mirrors/fl/FlexGen 点击查看 免费下载 DeepSpeed 在训练与推理时依赖一组 C/CUDA 扩展&#xff0…

作者头像 李华
网站建设 2026/9/25 8:00:36

自建CRM实战:DeskcommCRM部署、权限管理与数据安全指南

做销售管理的朋友&#xff0c;大概率都动过“自己搞一套CRM”的念头&#xff0c;尤其是当你发现市面上的免费CRM越用越别扭&#xff0c;收费CRM又贵得肉疼的时候。我团队之前就卡在这个点上&#xff0c;客户资料散在好几个人的微信和Excel里&#xff0c;月底统计全靠人工对表&a…

作者头像 李华