免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

爱奇艺实时数据架构演进:从Kafka到AutoMQ的存算分离实践

爱奇艺实时数据架构演进:从Kafka到AutoMQ的存算分离实践 2023年双十一大促前的压测周我记得很清楚凌晨三点多值班群里的告警一条接一条地顶上来线上某个Kafka集群的磁盘水位突破85%分区副本同步开始出现积压。这是爱奇艺实时流数据架构里最典型的一个场景——业务在增长流量在涨但支撑它们的Kafka集群越来越像一台开不动的老车。当时我们内部其实已经讨论了很久要不要动这套跑了几年的实时链路也就是把核心的Kafka架构迁到AutoMQ上。聊了半年真正推动落地的是那次压测里暴露出的扩容瓶颈Kafka集群的节点加不进去因为加节点要复制数据复制数据又要占带宽占带宽又把延迟打高延迟一高消费端积压积压一上来磁盘写放大更严重。这套恶性循环让团队下定决心做一次彻底的架构演进。这篇内容我打算把爱奇艺从Kafka到AutoMQ的整个演进过程完整复盘一遍包括我们为什么非换不可、AutoMQ的存算分离到底解决了什么、迁移的双跑与灰度怎么设计、踩了哪些坑、最终拿到了什么收益。如果你是做实时数据平台、消息中间件选型或者刚好也在评估是否要把自建的Kafka集群搬上云原生架构这篇应该能给你很具体的参考。1. 从Kafka集群的成长烦恼说起1.1 存储与计算强耦合扩容一个节点要连坐半天先交代一下背景。爱奇艺内部的实时数据链路覆盖的业务场景非常宽用户行为埋点采集、播放质量实时监控、内容推荐特征的实时计算、会员运营的实时营销触发还有大大小小的日志传输管道。这套体系里Kafka承担的是总线的角色几乎所有实时数据都要从它这里过一遍。高峰期集群日常处理的消息量级可以用每秒数百万条来估算再乘上多个业务域规模压力一下就出来了。传统Kafka的架构最核心的问题可以概括成一句话存储和计算强耦合。一个topic的所有分区必须固定落在某个broker节点的本地磁盘上消费者要消费某个分区就必须连接那个节点。这意味着扩容一个节点时Kafka要把一部分分区数据从旧节点复制到新节点整个过程中网络、磁盘、CPU都会被占掉一大块。我们测过一次一个中等规模的分区副本迁移能把同机架内其他broker的写入P99延迟抬高30%到50%。业务方那边立刻就有感知——消费延迟从正常的几十毫秒涨到几百毫秒甚至秒级。这种连坐效应在集群规模小的时候不明显毕竟节点少复制数据的带宽压力有限。但当我们把集群推到一个大几十节点的规模时任何一次扩缩容、任何一台机器下线都要面临数据搬家的时间成本。更麻烦的是Kafka在做副本迁移时如果同时撞上业务流量峰值整个集群的稳定性都会变得非常脆弱。1.2 分区迁移与Rebalance一次节点故障引发的共振另一个让我们头疼的点是Kafka的Rebalance机制。Kafka集群里如果有一台broker宕机或者一个分区leader切换该分区在新的leader上恢复服务前消费者组会触发Rebalance。这个过程本身没问题但在大集群里问题会被放大一个分区的leader切换引发消费者重连消费者重连导致一部分分区没有owner进而触发更多分区重新分配然后broker之间的同步流量又上升再引发下一轮抖动。我们内部管这种现象叫共振。实际案例我记得很清楚某次一台broker因为磁盘坏道被踢出集群结果后续两个小时里多个业务topic的消费延迟一路上涨运维同事在控制台上排查了很久才定位到是Rebalance风暴叠加了副本同步流量把集群的元数据请求也拖慢了。这种场景下Kafka自带的运维工具只能说够用但离好用差得很远。kafka有没有UI界面这类问题也是社区老话题Kafka原生确实只有命令行工具和JMX监控第三方UI基本是额外部署的排障链路越长稳定性的坑就越深。1.3 账单压力磁盘成本随留存周期指数上涨从成本视角看传统Kafka的存储模型更让我们焦虑。Kafka的数据要写本地磁盘为了保证可靠性默认要保存多副本我们生产环境用的是3副本。业务方经常要求topic数据保存7天甚至更长算一笔账就清楚了业务写入1TB数据落到集群里实际占用就是3TB磁盘再按7天留存算如果每天写入1TB集群里同时要存7天的数据总占用就是21TB。注意这21TB还是用3副本保护之后的量。磁盘价格虽然逐年下降但在我们这种规模下存储成本在实时链路总成本里的占比相当可观而且随着业务增长这部分是刚性上涨的。更不划算的是Kafka的存储利用率和访问频率其实很不匹配。大量的历史分区数据保存几天后根本没什么消费者去读了但它们仍然占据着本地磁盘、参与副本同步、消耗磁盘巡检和备份的资源。相当于你用高成本的本地盘养了一大堆冷数据。这个矛盾是Kafka架构天生带来的想解决只能换思路。1.4 弹性能力缺失大促峰值只能提前囤资源每年的大促和热门内容上线都是我们最紧张的时候。实时流量会突然冲到平时的三到五倍为了扛住峰值团队必须在活动前几周就开始准备预估峰值、临时扩容Kafka节点、业务活动结束后再缩容。听起来简单做起来全是泪。扩容节点要等数据迁移完成至少半天起步活动结束后缩容又要等数据搬走。很多时候流量峰值只持续两三个小时但资源你得按峰值备一整天。弹性能力缺失带来的是资源利用率长期偏低。平时集群负载可能只有三成但为了可能到来的峰值你得让集群维持一个冗余水位。这在业务负责人看来就是时时刻刻在烧钱。所以我们评估新的消息队列方案时弹性伸缩能力被列为和稳定性同等重要的指标。2. AutoMQ的破局逻辑存算分离如何打动我2.1 核心原理把数据搬家从扩容中拿掉AutoMQ打动我们的第一点是它的存算分离架构。它把Kafka里最重的两个负担——日志存储和分区元数据——从broker本地剥离出来交给云存储承担。具体来说写入的数据先落到本地的ESSD云盘做短时缓冲再异步上传到对象存储比如对象存储OSS/S3里做长期持久化。消费者读历史数据时优先命中本地缓存miss了才去远端拉取。这样带来的直接变化是一个分区的数据不再绑定在某个broker节点上节点只是计算层随时可以替换、增减。扩容一个节点时新节点不需要从旧节点复制全量数据它只需要接管一部分分区的计算负载元数据和存储都在远端读写路径可以立刻建立起来。分区迁移的时间从原来的小时级别直接降到秒级甚至毫秒级。我打个比方传统Kafka像一家每家分店都自带仓库的连锁餐厅新开一家店要把食材调拨过去路上还得防着食材变质AutoMQ更像中央厨房加冷链配送门店本身不需要囤货新开一家店接上配送线路就能开业。这个仓库和门店拆开的思路就是存算分离的本质。当然存算分离不是AutoMQ第一个提出的但它在工程实现上把细节做得比较扎实本地盘只承担热数据的读写对象存储承担全量数据的持久化这两层之间用WAL预写日志机制保证数据不丢。WAL先写本地确认成功后再返回客户端ack后台再异步刷到远端。这样既保证了低延迟又保证了数据可靠性。2.2 Kafka协议兼容这是决定性的入场券技术再先进如果让业务方改客户端代码这个项目在我们内部基本推不动。爱奇艺的业务团队用的Kafka客户端版本五花八门Java的、Go的、Python的还有基于librdkafka的C服务真要逐个改动辄几十个业务方配合周期会拖到以年计。AutoMQ在这方面的策略很务实——完全兼容Kafka协议。它对外提供的接入方式、Topic/Consumer Group的管理方式以及绝大多数客户端SDK的调用行为和Kafka保持一致。我们做一个POC概念验证时直接把业务方的Kafka客户端连接地址换成AutoMQ的broker地址什么都不改消息就能正常生产和消费。这一点太关键了它把迁移的成本从改造所有业务代码降到了改一个连接配置。当然完全兼容是相对的我们在POC阶段还是发现了一些细节差异。比如某些客户端版本对Kafka broker端返回的某些元数据字段比较敏感AutoMQ模拟这些字段时偶尔会有边界行为不一致。这部分我们在后面的灰度阶段做了适配但总体来看协议兼容带来的迁移成本降低是这次架构演进能快速推进的基础。2.3 Replica与故障转移的设计变化从ISR到WAL加远端存储传统Kafka的副本机制靠的是多个broker节点各自保存完整副本leader挂了从ISR同步副本集合里选一个follower提升上来。这套机制保证了高可用代价是副本之间的数据同步要占用网络和磁盘。AutoMQ的做法不一样——它把副本的概念往上抽象了一层数据先写到本地WAL再持久化到远端对象存储只要远端存储可靠即使某个broker节点整个宕机数据依然完整保存在对象存储里新节点接管分区时直接挂载远端的日志段即可。这意味着AutoMQ不需要在多个节点之间维护多份完整数据副本它的副本存储开销远低于Kafka。实际测试时我们模拟过节点故障AutoMQ的分区重新选举和服务恢复时间比Kafka要快很多因为新节点不需要从其他节点拉全量数据只需要从远端挂载数据再加上本地的热缓存预热整个恢复过程流畅很多。这里我不展开太多源码级细节但对运维团队来说这个变化带来的最直观好处是以前处理故障心里想的是会不会丢数据要多久才能补同步现在处理故障心里想的是换一台节点检查挂载状态。心智负担完全不是一个量级。3. 架构演进全过程从双跑验证到全量切换3.1 第一阶段影子集群与流量回放确定了要迁我们不能直接切毕竟这是承载核心业务数据的实时管道。我们设计的第一阶段是影子集群验证搭建一套和线上同规格的AutoMQ集群把线上Kafka的流量复制一份进去观察AutoMQ跑一段时间看延迟、吞吐、稳定性是否符合预期。具体做法很朴素我们用消费端把线上Kafka的Topic数据读出来再原样写入AutoMQ的对应Topic相当于做了一次流量回放。这个阶段不追求实时性只追求数据一致——我们对比两边集群的offset进度、消息条数、消息体内容确认AutoMQ写入和读取的数据没有丢失或篡改。这里要提醒一句流量回放时一定要注意幂等问题。我们的回放程序偶尔会因为消息重复投递导致写入重复数据对比时不要只看条数还要看具体的offset映射和消息内容哈希。我们内部写了一个小工具给每条消息计算CRC64然后对比两边相同offset的消息CRC能精确定位到任何不一致。影子集群我们跑了大概两周期间覆盖了业务峰值时段也人为做了多次故障演练。AutoMQ的表现基本稳定端到端延迟维持在几十毫秒级别这给了我们切换的信心。3.2 第二阶段按场景灰度迁移影子验证通过后进入灰度迁移阶段。我们没有搞一刀切而是把内部使用Kafka的场景分成三类按风险从低到高排序迁移日志传输类业务系统把日志实时发送到Kafka再由下游清洗入仓。这类数据允许少量延迟抖动丢了也能容忍作为第一批。异步消息类订单状态变更、通知推送等业务消息。要求不能丢但对延迟的容忍度在秒级作为第二批。核心实时链路类用户行为埋点、推荐特征、播放质量监控。这是实时性要求最严的场景延迟要求百毫秒内作为最后一批。灰度迁移的技术方案我们采用了双写消费者切换的组合。业务方在生产端写消息时同时写入Kafka和AutoMQ保证两边数据都有然后我们把消费者从Kafka逐步切到AutoMQ上观察消费端是否正常指标是否稳定。双写期间如果AutoMQ出问题消费者随时可以切回Kafka回退成本很低。这里有一个细节值得分享双写不是简单的写两份相同数据两个消息队列的写入延迟可能不同导致消费端拉取时出现乱序。我们的做法是在消息体里带上生产时间戳和消息序号消费端做一次轻量的排序校验。这个保护机制在切换初期特别有用能有效降低数据好像没问题但业务逻辑却异常的排查难度。3.3 第三阶段大促峰值验证与全量切换灰度迁移期间刚好赶上一次常规的大型促销活动我们索性把AutoMQ作为主力集群来扛这次峰值。这是最真实的检验流量是真实的业务压力是真实的如果AutoMQ在这个场景下能稳定扛住我们就有底气做全量切换。那次峰值AutoMQ的表现让我们很惊喜。因为存算分离扩容操作变成了加节点这么简单新节点起来后几乎立即投入服务不再需要等数据迁移。我们在峰值前半小时临时加了两台节点整个操作花了不到十分钟这在传统Kafka集群上是不可想象的——以前加两个节点光数据迁移就得等四五个小时。峰值过后我们完成了剩余场景的切换并逐步把Kafka集群缩容。全量切换后团队又运行了两周确认无异常才正式停掉大部分Kafka节点只保留了少量节点作为紧急回退通道。4. 演进落地后的实测数据与关键收益4.1 延迟指标P99稳定在百毫秒以内先看大家最关心的延迟。迁移完成后我们重点盯了三种指标生产端写入延迟、消费端端到端延迟、大消息超过1MB的单条消息的处理耗时。实测下来AutoMQ在正常负载下的写入延迟和Kafka基本持平甚至因为不再受磁盘抖动影响P99延迟更稳定。端到端延迟上AutoMQ在本地缓存命中率高的情况下大部分场景能做到百毫秒以内只有冷数据回放时因为要访问对象存储延迟会上升到几百毫秒。这也是我们后续重点优化的方向之一后面会讲到。之前群里总有同事问kafka消息延迟高怎么排查我们在Kafka时代也做过很多这类问题定位多半和磁盘IO、分区不均衡、消费者处理慢相关。AutoMQ时代延迟高的原因更集中在对象存储的读放大上排查思路有很大不同。4.2 成本账单存储开销下降超七成成本是我们这次演进最直观的收益。AutoMQ的存储模型下热数据在本地盘只保留一小段时间全量数据在对象存储。对象存储的单GB成本远低于我们自建Kafka所用的本地SSD再叠加AutoMQ不需要3副本的本地存储整体存储成本下降非常明显。我们不严谨地算过一笔账以我们一个日均写入2TB、留存7天的核心集群为例——Kafka时代需要准备约42TB的本地磁盘容量2TB×7天×3副本数据放本地盘AutoMQ时代本地盘只需要缓冲少量的热数据比如几百GB长尾数据全部落到对象存储。单纯从存储成本看下降超过七成。如果算上因弹性能力带来的计算资源节省不需要长期保持冗余水位整体TCO下降比例还会更高。这里必须说明成本测算不能只看存储单价。AutoMQ的broker节点对CPU和内存的要求仍然存在如果你把Kafka集群已经压得很满那么切到AutoMQ后计算资源并不会减少太多。但从我们的实际账单看存储成本的下降已经足以覆盖这部分。4.3 运维体验从救火到规划架构演进带来的另一个好处是运维团队的工作方式发生了根本变化。以前我们最怕的就是节点磁盘快满了和分区不平衡这两类问题占用了大量排障时间。AutoMQ因为数据都在远端存储本地盘更像是一个缓存盘满了就清、坏了就换不会出现某个分区数据卡死在坏盘上的尴尬。另外AutoMQ提供了一套可视化的控制台支持查看集群状态、Topic流量、消费组延迟等指标。这弥补了Kafka原生缺乏UI界面的短板也让我们的值班同学排查问题时更快定位到瓶颈。我们内部再也没有出现过凌晨三点靠命令行一条条翻日志的场景。关于kafka集群安装和kafka kafka教程这类话题我也想多说一句如果你还在自建Kafka集群一定要重视自动化部署和监控平台的建设因为Kafka的运维复杂度是随规模上升的。我们的经验是自建Kafka的边际运维成本一直在涨这也是促使我们拥抱云原生架构的原因之一。5. 踩过的坑和值得记住的经验5.1 坑一对象存储的读放大与热分区缓存AutoMQ把长期数据放到对象存储后一个新的问题出现了某些高热度分区的数据在消费者追尾读取时会频繁访问远端对象存储造成读放大。尤其是消费者重放场景比如从某个旧offset重新消费大量的远距离读会让对象存储的访问延迟成为瓶颈。我们踩了一次很深的坑某个业务方做数据修复从三天前的offset重新消费一个每天写入量巨大的topic结果AutoMQ集群的broker CPU突然飙升因为每个节点都在从对象存储拉取大量历史数据然后把数据在本地做索引重建。我们后来做了两个优化一是为高频消费的数据增加本地缓存策略二是限制单租户的远距离读取并发把大流量回放任务拆分到低峰期执行。如果你也要上AutoMQ建议提前考虑好热数据缓存命中率这个指标最好在压测阶段就覆盖消费者追尾读场景。5.2 坑二小分区多Topic场景的碎片化爱奇艺内部有很多业务方习惯一个业务一个Topic导致集群里出现大量的小Topic每个Topic的分区数很少但总量很多。在AutoMQ架构下这些小Topic的数据会零散地存到对象存储产生大量小文件碎片。对象存储对小文件的读写性能远低于大文件这会导致一定程度的性能劣化。我们的运维建议是迁移前和业务方一起做一次Topic治理把强相关的业务合并到一个Topic里用分区key区分业务类型。这听起来像脏活累活但能有效提升AutoMQ的性能表现。如果你在评估阶段发现自己的集群大量是小Topic一定要先做治理再迁移不然后期会很难受。5.3 坑三客户端版本与协议兼容的边界差异虽然AutoMQ兼容Kafka协议但我们在灰度阶段还是遇到了两个和客户端版本兼容相关的边界问题。一个是某个老版本的Java客户端在获取topic元数据时使用了过时的协议字段AutoMQ这边返回的更精确的值反而让客户端报错另一个是某些客户端在Consumer Group Rebalance时的行为和原生Kafka存在细微差异导致有短暂的消费停顿。这两个问题都不难解决——升级客户端版本或者调整AutoMQ端的兼容参数——但它们提醒我们兼容性不能只看能不能连上要覆盖到元数据交互、Rebalance行为、配额管理这些细节。建议你在做POC时把线上主流客户端版本都列出来逐个做连接和压力测试不要只测最新版本。5.4 给后来者的迁移Checklist最后整理一份我们实际操作沉淀下来的迁移Checklist方便团队直接抄作业阶段检查项注意事项选型评估协议兼容性测试覆盖线上所有客户端版本不只测新版选型评估存算分离架构理解区分热数据本地缓存和冷数据远端读影子验证流量回放的一致性用消息CRC校验不只对比条数影子验证故障演练模拟节点宕机、机房异常等场景灰度迁移场景分级排序先日志类再业务消息类最后核心链路灰度迁移双写机制设计消息体打时间戳和序号防乱序全量切换回退预案保留少量Kafka节点作为紧急通道运营阶段监控指标建设重点看对象存储读放大和缓存命中率如果你也想从Kafka迁到AutoMQ我建议你先干三件事第一梳理自己的Topic和流量特征弄明白哪些是热点数据哪些是长尾数据第二建一套流量回放和一致性比对的工具这是整个迁移的安全网第三别追求一步到位把迁移拆成多个批次每一批都要有独立的回退方案。我在实际推进这个项目的过程中最深的一个体会是架构演进这件事技术选型本身只占三成精力剩下七成都在做业务协调、风险控制和运维体系重建。Kafka不是不能用了而是当我们的实时流数据场景发展到一定规模传统架构带来的存储成本、扩容效率和运维负担会逐步压过它的成熟稳定优势。AutoMQ这种云原生存算分离的思路正好切中了我们最痛的那几个点。如果你所在的团队也面临着类似的存储成本压力和弹性扩容诉求不妨认真评估一下这条技术路线也许它能帮你把实时数据架构从一个沉重的包袱变成业务增长真正的助力。
返回列表