
AutoMQ 无盘 Kafka 架构解读基于 S3 的共享存储、秒级扩缩容与 10x 成本优化实践【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automqAutoMQ 是 Apache Kafka 的一个云原生分支fork它用对象存储S3 或任何 S3 兼容存储如 MinIO重构了 Kafka 的存储层将经典的 shared-nothing 架构转变为 shared storage 架构从而实现无状态 Broker、秒级扩缩容与跨 AZ 流量成本归零。本文以仓库 README.md 为主体骨架结合 docker/docker-compose.yaml、s3stream 与 AutoMQConfig.java 等源码带你完成本地单节点/三节点集群部署并深入理解其架构组件与关键配置项的底层实现。一、AutoMQ 是什么从 Kafka 的三大痛点说起AutoMQ 定位为运行在 S3或任何 S3 兼容存储之上的无状态 Kafka 替代方案。它之所以被设计出来是为了解决 Apache Kafka 在云环境下面临的两个核心问题弹性伸缩困难Kafka Broker 是有状态的扩缩容需要移动数据即使只是在 Broker 之间 reassign 分区也是一个复杂过程云上托管成本高昂EBS 存储费用、跨 AZ 流量费用加上 Kafka 有限的伸缩能力导致的严重过度预置over-provisioning让 Kafka 集群的账单居高不下。AutoMQ 的思路是用共享存储替换本地磁盘把数据层下沉到对象存储让计算层Broker彻底无状态化。其核心亮点可以归纳为成本效益Cost effective首个真正云原生的流式存储系统面向云环境做成本与效率的极致优化高可靠High Reliability依托对象存储服务实现零 RPO、秒级 RTO以及 99.999999999% 的数据持久性零跨 AZ 流量Zero Cross-AZ Traffic以云对象存储作为首选存储方案消除 AWS/GCP 上的跨 AZ 流量费用传统 Kafka 集群中超过 80% 的成本来自生产者、消费者与副本同步侧的跨 AZ 流量Serverless 化监控集群指标自动扩缩容、计算层Broker无状态可在秒级完成伸缩、以对象存储为主存储消除了容量焦虑免运维Manage-less内置 Auto Balancer 组件自动在 Broker 之间调度分区与网络流量免去手动分区 reassign高性能High Performance通过预取pre-fetching、批处理与并行技术最大化对象存储的能力内置指标导出Built-in Metrics Export原生导出 Prometheus 与 OpenTelemetry 指标同时支持 push 和 pull 两种模式告别低效的 JMX100% Kafka 兼容完整兼容 Apache Kafka提供全部功能的同时拥有更优的成本与运维效率。二、快速开始Docker Compose 本地部署 AutoMQ仓库在 docker 目录下提供了两套 Docker Compose 编排文件用于快速评估 AutoMQ 能力。2.1 前置条件在本地运行 AutoMQ 之前请确保环境满足以下要求Docker 版本 20.x 或更高Docker Compose v2为 Docker 分配至少 4 GB 内存系统上 9092 和 9000 端口可用注意生产级 AutoMQ 集群的部署是有挑战的本快速开始仅用于评估 AutoMQ 特性不适合生产环境使用。2.2 单节点快速启动docker/docker-compose.yaml 提供了一个简单的单节点配置适合快速评估与开发curl -O https://raw.githubusercontent.com/AutoMQ/automq/refs/tags/1.5.5/docker/docker-compose.yaml docker compose -f docker-compose.yaml up -d该配置包含一个同时承担 controller 和 broker 角色的 AutoMQ 单节点以及一个用于提供 S3 存储的 MinIO。所有服务运行在名为automq_net的 Docker bridge 网络中你可以在该网络中启动一个 Kafka 生产者来测试 AutoMQdocker run --network automq_net automqinc/automq:latest /bin/bash -c \ /opt/automq/kafka/bin/kafka-producer-perf-test.sh --topic test-topic --num-records1024000 --throughput 5120 --record-size 1024 \ --producer-props bootstrap.serversserver1:9092 linger.ms100 batch.size524288 buffer.memory134217728 max.request.size67108864测试完成后可用以下命令销毁环境docker compose -f docker-compose.yaml down2.3 单节点编排文件逐段解析让我们对照 docker/docker-compose.yaml 理解这个最小 AutoMQ 集群由哪些部分构成MinIO 服务S3 存储层minio: container_name: minio image: minio/minio:RELEASE.2025-05-24T17-08-30Z environment: - MINIO_ROOT_USERminioadmin - MINIO_ROOT_PASSWORDminioadmin - MINIO_DOMAINminio ports: - 9000:9000 # MinIO API - 9001:9001 # MinIO ConsoleMinIO 扮演 AutoMQ 的对象存储后端9000 端口暴露 S3 API9001 端口是其 Web Console。健康检查通过curl http://minio:9000/minio/health/live完成。mc 服务自动初始化桶mc: image: minio/mc:RELEASE.2025-05-21T01-59-54Z depends_on: minio: condition: service_healthy entrypoint: /bin/sh -c until (/usr/bin/mc alias set minio http://minio:9000 minioadmin minioadmin) do echo ...waiting... sleep 1; done; /usr/bin/mc rm -r --force minio/automq-data; /usr/bin/mc rm -r --force minio/automq-ops; /usr/bin/mc mb minio/automq-data; /usr/bin/mc mb minio/automq-ops; /usr/bin/mc anonymous set public minio/automq-data; /usr/bin/mc anonymous set public minio/automq-ops; tail -f /dev/null mc 容器在 MinIO 健康后自动创建两个桶automq-data数据桶与automq-ops运维/元数据桶并将其设为 public 匿名可读供 AutoMQ 读写。AutoMQ 单节点controller broker 合并角色server1: container_name: automq-single-server image: automqinc/automq:1.6.0 stop_grace_period: 1m environment: - KAFKA_S3_ACCESS_KEYminioadmin - KAFKA_S3_SECRET_KEYminioadmin - KAFKA_HEAP_OPTS-Xms1g -Xmx4g -XX:MetaspaceSize96m -XX:MaxDirectMemorySize1G # Replace CLUSTER_ID with a unique base64 UUID using bin/kafka-storage.sh random-uuid - CLUSTER_ID3D4fXN-yS1-vsQ8aJ_q4Mg command: - bash - -c - | /opt/automq/kafka/bin/kafka-server-start.sh \ /opt/automq/kafka/config/kraft/server.properties \ --override cluster.id$$CLUSTER_ID \ --override node.id0 \ --override controller.quorum.voters0server1:9093 \ --override controller.quorum.bootstrap.serversserver1:9093 \ --override advertised.listenersPLAINTEXT://server1:9092 \ --override s3.data.buckets0s3://automq-data?regionus-east-1endpointhttp://minio:9000pathStyletrue \ --override s3.ops.buckets1s3://automq-ops?regionus-east-1endpointhttp://minio:9000pathStyletrue \ --override s3.wal.path0s3://automq-data?regionus-east-1endpointhttp://minio:9000pathStyletrue这段命令揭示了 AutoMQ 使用 KRaft 模式启动并通过--override注入三类关键配置KRaft 元数据cluster.id用bin/kafka-storage.sh random-uuid生成唯一 base64 UUID、node.id、controller.quorum.voters与controller.quorum.bootstrap.servers网络监听advertised.listenersPLAINTEXT://server1:9092S3 存储三要素s3.data.buckets数据桶、s3.ops.buckets运维桶、s3.wal.pathWAL 路径均指向 MinIO 的 S3 endpoint并使用pathStyletrue兼容 MinIO 的路径式访问。2.4 三节点集群部署docker/docker-compose-cluster.yaml 提供了更复杂的三节点 AutoMQ 集群配置适合测试 AutoMQ 的集群特性运行方式与单节点一致。该文件通过 YAML 锚点x-common-variables: common-env复用公共环境变量S3 访问密钥、堆参数、CLUSTER_ID三个节点server1/server2/server3的关键差异仅在--override node.id0 # server1 --override node.id1 # server2 --override node.id2 # server3 --override controller.quorum.voters0server1:9093,1server2:9093,2server3:9093 --override controller.quorum.bootstrap.serversserver1:9093,server2:9093,server3:9093 --override advertised.listenersPLAINTEXT://server1:9092 # 各节点各自的主机名三节点共享同一个s3.data.buckets/s3.ops.buckets/s3.wal.path——这正是共享存储架构的直接体现所有 Broker 读写同一份对象存储中的数据Broker 本身不持有数据因此任何一个节点都可以随时加入或退出集群而无需进行数据搬迁。此外仓库还提供了更多部署选项详见 README.mdDocker 多节点测试集群、Linux 5 节点集群、Kubernetes 部署以及 AWS / 阿里云 Marketplace 的两周免费试用入口。三、核心架构从 shared-nothing 到 shared storageAutoMQ 是 Apache Kafka 的开源分支引入基于对象存储的新存储引擎把经典的 shared-nothing 架构改造成 shared storage 架构AutoMQ 与 Kafka 的架构本质区别在于存储层AutoMQ 利用对象存储实现无状态 Broker 架构。它由以下关键组件构成组件职责S3 Storage Adapter适配层重新实现了 Kafka 的UnifiedLog、LocalLog与LogSegment类将日志创建在 S3 上而非本地磁盘如需要仍可支持传统本地磁盘存储S3Stream共享流式存储库封装了 WAL写前日志与对象存储等多个存储模块WAL 面向频繁写入与低 IOPS 场景做了专门优化以降低 S3 API 成本为提升读性能引入 LogCache 与 BlockCacheAuto Balancer自动平衡 Broker 之间的流量与分区免去手动 reassign不同于 Kafka这一内置能力替代了对 Cruise Control 的依赖Rack-aware RouterKafka 长期面临 AWS/GCP 跨 AZ 流量费用问题共享存储架构通过 rack-aware 路由器为不同 AZ 的客户端提供特定的分区元数据数据经由对象存储交换从而规避跨 AZ 费用3.1 S3Storage源码视角的共享存储实现在 s3stream/src/main/java/com/automq/stream/s3/S3Storage.java约 997 行中可以看到共享存储的核心类S3Storage。从其依赖可以看出读写链路的完整组成写入链路WriteAheadLogs3stream/src/main/java/com/automq/stream/s3/wal负责接收高频写入LogCachecom.automq.stream.s3.cache.LogCache缓存尚未上传或已上传未淘汰的数据ObjectManager与ObjectWriter负责将数据对象化写入对象存储读取链路S3BlockCachecom.automq.stream.s3.cache.S3BlockCache缓存从对象存储读取的冷数据ObjectReader与CompositeObjectReader负责读取对象内容LocalStreamRangeIndexCache提供流范围索引缓存故障恢复Failover与StorageFailureHandler处理存储故障时的接管逻辑WALRecovery负责从 WAL 恢复数据。这种设计使得写先落 WAL 快速确认、再异步批量上传到 S3从而以少量 S3 API 调用支撑高频写入——这正是 AutoMQ 声称可以将 Apache Kafka 云上账单削减 90%的实现基础具体基准数据可参考仓库 README 中引用的成本对比报告。3.2 S3Stream 配置项核心参数与默认值s3stream/src/main/java/com/automq/stream/s3/Config.java 定义了 S3Stream 层的核心参数及其默认值配置项默认值含义walConfig0file:///tmp/s3stream_walWAL 路径配置生产环境通常指向 S3walCacheSize200 MBWAL 缓存大小FIFO 队列容纳未上传及已上传未淘汰的数据walUploadThreshold100 MBWAL 触发上传的阈值walUploadIntervalMs-1不按时间上传WAL 按时间间隔上传的周期streamSplitSize16 MB上传增量 WAL 或压缩 stream set object 时的切分阈值objectBlockSize1 MBS3 对象压缩块大小阈值objectPartSize16 MBS3 对象分片上传multi-part upload的分片大小阈值blockCacheSize100 MBBlockCache 大小缓存从对象存储读取的冷数据streamObjectCompactionIntervalMinutes60 分钟Stream object 压缩周期越大 API 调用成本越低但元数据规模越大streamObjectCompactionMaxSizeBytes10 GBStream object 压缩允许合成的最大对象大小networkBaselineBandwidth1 GB/s对象存储请求的总可用带宽防止压缩与追赶读垄断正常读写流量objectRetentionTimeInSecond600 秒10 分钟标记删除的 S3 对象保留时间3.3 Broker 侧配置AutoMQConfig 的关键参数Broker 侧的 S3 相关配置定义在 core/src/main/java/kafka/automq/AutoMQConfig.java与 docker-compose 中的--override一一对应s3.data.buckets数据桶地址格式0s3://$bucket?region$region可附加endpoint与pathStyle参数s3.ops.buckets运维桶地址格式与 data buckets 相同s3.wal.pathWAL 路径格式0s3://$bucket?region$region[batchInterval250][maxBytesInBatch8388608]即支持配置批处理间隔默认 250ms与单批最大字节数默认 8 MBs3.wal.cache.sizeWAL 缓存大小作为 FIFO 队列容纳尚未上传到对象存储的数据以及已上传但尚未从缓存淘汰的数据s3.stream.object.split.size上传增量 WAL 或压缩 stream set object 时的切分阈值s3.object.block.sizeS3 对象压缩块大小阈值s3.object.part.sizeS3 对象分片上传的分片大小阈值s3.block.cache.sizeBlockCache 大小用于缓存从对象存储读取的冷数据s3.stream.allocator.policyS3 流内存分配策略支持 DIRECT 内存模式使用 DIRECT 时需同步调整-Xmx与-XX:MaxDirectMemorySize可通过环境变量KAFKA_HEAP_OPTS设置s3.stream.object.compaction.interval.minutes默认 60 分钟与s3.stream.object.compaction.max.size.bytes默认 10 GBStream object 压缩的周期与上限二者共同权衡 API 调用成本与元数据规模s3.stream.set.object.compaction.interval.minutes默认 5 分钟Stream set object 压缩周期值越小元数据规模越小、数据越早可压缩但对象经历的压缩次数增多s3.stream.set.object.compaction.stream.split.size默认 8 MB压缩过程中单流数据超过该阈值则直接拆分写入单个 stream objects3.stream.set.object.compaction.force.split.minutes默认 120 分钟stream set object 压缩的强制拆分周期s3.network.baseline.bandwidth默认 1 GB/s对象存储请求总可用带宽用于防止 stream set object 压缩与追赶读catch-up read垄断正常读写流量s3.object.delete.retention.minutes默认 10 分钟标记删除对象在 S3 中的保留时间。这些参数共同构成了一套以 API 调用成本换存储规模 / 以存储规模换 API 调用成本的可调杠杆是 AutoMQ 实现低成本的核心旋钮。四、Auto Balancer免去手动 reassign 的内置负载均衡AutoMQ 的内置 Auto Balancer 组件位于 core/src/main/java/kafka/autobalancer其核心类AutoBalancerManager负责自动调度分区与网络流量。从源码结构看它由以下模块组成指标采集AutoBalancerMetricsReporter周期性上报 Broker 与分区维度的负载指标见 metricsreporter指标模型定义在TopicPartitionMetrics/BrokerMetrics中异常检测AnomalyDetectorImpl基于指标快照检测负载不均衡的异常目标定义Goal接口及其实现如NetworkOutUsageDistributionGoal、NetworkInUsageDistributionGoal定义了均衡目标——从实现上看Auto Balancer 同时关注网络流入与流出的分布执行器ControllerActionExecutorService在 Controller 侧执行分区迁移等动作。对运维而言这意味着 AutoMQ 集群不需要像 Kafka 那样手动执行分区 reassign 或额外部署 Cruise Control负载均衡由系统自动完成——这正是 README 中 Manage-less 体验的落地点。五、Rack-aware Router从架构层面消除跨 AZ 流量费跨 AZ 流量费是云上 Kafka 成本的大头。README 明确指出在传统 Kafka 集群中超过 80% 的成本来自跨 AZ 流量包括生产者、消费者与副本同步三侧。AutoMQ 的解法是架构性的由于所有数据都在对象存储天然多 AZ 冗余中Broker 不再需要跨 AZ 同步副本。Rack-aware Router 的作用是为不同 AZ 的客户端提供特定的分区元数据——客户端被路由到与其同 AZ 的 Broker 进行读写数据本身通过共享的对象存储交换从而把跨 AZ 流量费用降到零。这一点也呼应了 README 中 Zero Cross-AZ Traffic 的承诺在 AWS 与 GCP 上使用 AutoMQ不再产生跨 AZ 流量账单。六、最新特性Table Topic——流表一体统一流式与数据分析Table Topic 是 AutoMQ 的新特性将 stream 与 table 功能结合统一流式处理与数据分析。当前它支持 Apache Iceberg并集成了 AWS Glue、HMS、Rest catalog 等 catalog 服务同时原生支持 S3 Tables——AWS 在 2024 re:Invent 上发布的新产品。从仓库代码看Table Topic 的实现位于 core/src/main/java/kafka/automq/table其中TableManager.java 负责表的管理生命周期TopicPartitionsWorker.java 与 WorkerConfig.java 定义了执行写入/转换任务的工作线程及其配置。在 docker/table_topic 目录下还提供了基于 Spark Iceberg 的入门 NotebookTable Topic - Getting Started.ipynb与配套 Dockerfile方便在本地快速体验 Table Topic 的流表一体能力。七、参与社区与许可提问与报 bug通过 GitHub Issues交流讨论加入 Slack 频道或微信群二维码见 docs/images/automq-wechat.png贡献代码先阅读 CODE_OF_CONDUCT.md 与 CONTRIBUTING_GUIDE.md仓库维护了 good first issues 列表帮助新贡献者入门许可证AutoMQ 采用 Apache 2.0 许可证详见 LICENSE商标声明Apache®、Apache Kafka®、Kafka®、Apache Iceberg®、Iceberg® 及相关开源项目名称为 Apache Software Foundation 的商标。小结通过本文你可以看到AutoMQ 的价值主张并非停留在口号层面而是有完整的代码与配置支撑docker/docker-compose.yaml展示了 5 分钟即可体验的无盘 KafkaS3Storage与Config.java揭示了 WAL LogCache/BlockCache 对象存储的分层读写设计AutoMQConfig中的一系列s3.*参数给出了成本/性能的可调旋钮kafka/autobalancer与 rack-aware 路由则在架构层面兑现了免运维与零跨 AZ 流量的承诺。对于正在评估从 Kafka 迁移到云原生流式存储的团队可以按仓库 docker 目录下的编排文件先行验证再结合 config/kraft 中的配置模板规划生产部署。【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automq创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考