ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka+MongoDB:大数据文档存储的稳定搭档与实战指南

Kafka+MongoDB:大数据文档存储的稳定搭档与实战指南 做大数据相关项目这几年接手的文档存储需求一茬接一茬从埋点日志到业务事件从接口回调到网约车轨迹几乎都是半结构化的 JSON。早期有人图省事直接塞 MySQL结果字段一对不齐就要改表结构数据量上来以后查询和写入都得跪着调也有人一上来就扔 Elasticsearch查询是快了但作为“事实存储”总感觉少了点东西。折腾过一圈之后我固定下来的一套组合拳就是Kafka 负责接入和缓冲MongoDB 负责落地与查询。两个系统各管一段中间用消费管道串起来这就是这篇要讲的“Kafka 与 MongoDB 在大数据文档存储中的协作”。这篇内容适合正在做 Kafka 数据接入、MongoDB 存储选型或者打算搭一套实时文档管道的人看。我会直接讲架构思路、配置参数、代码细节和踩过的坑不写那种复制官方文档的“教程”。1. 为什么这套组合适合大数据文档存储1.1 两个角色一条流水线加一个仓库我在给新人解释这套架构时最喜欢用的例子是“流水线加仓库”。Kafka 就是流水线它不关心箱子里面装的是什么只负责把箱子从 A 点运到 B 点速度快、吞吐高、能扛住瞬间大量货物涌入。MongoDB 是仓库箱子到了以后按品类、按货架码放整齐随时可以查、可以改、可以按条件翻出来。流水线和仓库之间没有硬耦合。生产端业务系统不用等存储层响应消息丢进 Kafka 就算“交接”了存储层也不用跟生产端抢资源消费端按自己的节奏从 Kafka 拉数据批量写进 MongoDB。这个“异步解耦”恰恰是大数据文档存储里最值钱的一点写入峰值来了Kafka 先把数据削平MongoDB 不会被打爆。拿我之前做过的一个用户行为采集项目来说周末大促时段每秒写入峰值能冲到几千条平时只有几百。如果直连 MongoDB 写入连接池、磁盘 IO 都会被瞬时打满应用层还得做各种降级。加上 Kafka 之后生产端只管往 topic 塞消息MongoDB 消费端按固定线程数、固定批量大小慢慢落库整个链路非常稳。这里的本质是用 Kafka 的“吸收能力”换存储层的“稳态吞吐”。1.2 文档数据的脾气正好 MongoDB 接得住文档型数据和关系型数据的最大区别在于结构不固定、嵌套多、字段会变。一条订单数据上午可能只有金额和商品 ID下午就可能多了优惠券明细再隔两天又冒出个物流信息。你去设计 MySQL 表结构要么冗余一堆空字段要么频繁加列而 MongoDB 的集合没有强制 schema每条文档按自身结构存储天然贴合这种“随业务变化”的数据。字段嵌套的查询MongoDB 也做得不错。比如文档里有个items数组里面每个元素又有商品、数量、价格要查“某一笔订单里某个商品的数量”一条查询就能精确匹配到数组元素这在关系型里要拆三张表再 JOIN复杂度和性能都上去了。底层存储引擎也值得一说。MongoDB 使用的WiredTiger 存储引擎文档默认采用压缩存储像我们那种大量重复 key 的日志类 JSON压缩比非常可观。相比纯文本文件或关系型表的存储方式磁盘占用少扫描 IO 也随之降低。1.3 我靠着这对组合解决了哪些具体痛点把 Kafka 和 MongoDB 串在一起我实际解决过下面这四类问题削峰填谷写入量有突发性通过 Kafka 缓冲MongoDB 稳定消费集群不会因为瞬时高峰被压垮。结构多变上游接口 JSON 字段每天都在变MongoDB 不用改表结构直接存查询时按需取字段。链路解耦数据生产方不用关心存储细节存储层挂了也不会阻塞业务写入消息先攒在 Kafka恢复后继续消费。水平扩展Kafka 靠加分区扩展吞吐MongoDB 靠分片扩展存储和查询两边都是横向伸缩不会撞到单机天花板。这些痛点几乎是大数据文档存储的全部核心诉求。只要你的场景里占了两条以上这套组合就值得认真考虑。2. 协作模式与架构设计动手前先想清楚2.1 三种协作方式我建议这么选Kafka 和 MongoDB 协作从实现角度无非三条路官方 Kafka Connect 插件、自研消费端、Change Streams 反向同步。我把它们的差异整理成了一张表。协作方式优点缺点适用场景Kafka Connect Mongo Sink Connector配置即用、免代码、支持批量与幂等定制逻辑受限排错要熟悉 Connect 框架标准 JSON 入库存字段不用复杂加工自研消费者 bulkWrite灵活可控可以任意清洗、过滤、路由要自己处理提交位移、重试、幂等需要对消息加工、拆分、合并或者写入前要做业务校验Mongo Change Streams → 写回 Kafka能捕获 Mongo 变更事件满足双向协作对 Mongo 负载有影响链路更长需要把存储层变更同步给其他系统大多数项目我会优先推荐第一种性能数据同步稳定伸缩直接调 task 数。但如果你的消息要经过字段映射、敏感信息过滤或者要按不同业务域写到不同集合那建议走第二种代码里什么都能做。Change Streams 那条路适合特定的同步场景不是本文重点。2.2 Topic 与分区设计顺序、并发都藏在这里数据要在 Kafka 里走一遭最先决定的是Topic 怎么划分。我的习惯是按业务域建 topic比如order-event、user-log、pay-notify不要一个 topic 装天下。原因很简单topic 是消费端隔离和扩展的基本单位混在一起后清洗逻辑、消费并发、错误处理都会缠成一团。分区数决定了消费并行度上限。一个分区同一时刻只能被同一个消费组里的一个消费者线程消费所以想让落库更快分区数不能太少。我一般按“预估峰值写入速率/单消费者处理速率”来粗算。比如每秒需要消费 5000 条单个消费者批量写 MongoDB 实测每秒能处理 1000 条那至少需要 5 个分区。当然落库能力和批量大小、文档大小、MongoDB 集群配置都有关系上线后还要实测调整。分区的另一个作用是保存顺序。Kafka 只能保证分区内消息有序不是全局有序。所以如果你要求同一订单的变更消息按时间顺序落库发送时就要用订单号做 key让相同 key 的消息进同一个分区。没有顺序要求的日志类数据key 可以不用或者随机化这样分区负载更均匀。这里有个容易犯的错把分区数一开始就设得很大以为能留足扩展余量。实际上分区数过大会增加 rebalance 耗时和文件句柄开销。更合理的做法是先用计算值后续真遇到瓶颈再按 key 重新设计消息路由用新 topic 推进。2.3 文档模型设计这是落库前最重要的一次决策MongoDB 的文档模型必须在管道还没跑起来时就设计清楚。动手前问自己几个问题_id 用什么消息里如果有天然业务唯一键比如订单号、设备 ID可以把它直接作为 _id 字符串这样写入时天然去重upsert 也很方便。如果没有唯一键再用 MongoDB 自动生成的 ObjectId。字段嵌套深度MongoDB 单文档有 16MB 大小限制普通日志到不了这个值但嵌套无限叠加会带来查询和更新成本。一般两层到三层就到头了再深就该拆文档或改结构。要建哪些索引按查询习惯建索引别等数据量上来再补。比如经常按userId createTime查就建联合索引日志类数据按时间清理就加 TTL 索引让 MongoDB 自动删过期文档。要不要分片单机超过几百 GB 或者写入 QPS 持续走高就要计划分片键。分片键要选基数大、分布均匀的字段比如userId千万别用只有几个取值范围的字段。这些决策直接影响后续查询性能和运维成本。我见过太多项目topic 建好了、消费端跑起来了结果 Mongo 集合没有索引一张几千万文档的集合每次查询都全表扫描只能停下来补索引导致链路阻塞。所以文档模型永远是先行的Kafka 那头再快也救不了存储层设计上的懒。3. 环境搭建与关键参数解析3.1 Kafka 集群的部署选择与配置Kafka 现在主推 3.x 版本最大的变化是KRaft 模式可以不用再依赖 ZooKeeper部署更简单元数据处理也更稳定。本地实验或中小规模集群我建议直接上 KRaft如果团队里还是老一套 ZooKeeper 集群也不是不能用只是新项目没必要再用旧模式。三节点是最常见的生产集群规模。每台节点建议 8 核 CPU 以上、16GB 内存起步磁盘用 SSD日志目录单独挂盘。Kafka 本身对磁盘 IO 敏感机械盘在高峰期很容易成为瓶颈。单机实验环境可以用 KRaft 快速起一个节点大致步骤是# 下载 Kafka 3.x 并解压后先格式化存储目录 bin/kafka-storage.sh random-uuid /tmp/kafka-uuid bin/kafka-storage.sh format -t $(cat /tmp/kafka-uuid) -c config/kraft/server.properties # 然后启动节点 bin/kafka-server-start.sh config/kraft/server.properties启动后创建 topic 用kafka-topics.sh --create --topic test-doc --partitions 3 --replication-factor 1验证生产和消费可以分别用kafka-console-producer.sh和kafka-console-consumer.sh。这里有两个参数要留意log.retention.hours决定消息保留时间文档存储场景如果还有离线分析要做建议保留 24 到 72 小时num.partitions是默认分区数只影响未显式指定分区数的 topic别指望它代替业务设计。3.2 Linux 下 MongoDB 安装与副本集初始化MongoDB 的安装社区版直接从官方源装最省事。以 Ubuntu 为例先导入公钥、添加官方源然后 apt 安装。很多人在这一步失败原因基本是网络源不通或者公钥过期解决办法是换用镜像源并更新公钥。安装完成后第一件事就是把 MongoDB 配成副本集哪怕当前只有一台节点。因为后面无论是用 MongoDB 官方 Kafka Connector还是想用 Change Streams都需要副本集支持。单节点也能初始化配置文件里加一行replication.replSetName: rs0重启后执行rs.initiate()看到ok: 1就说明单节点副本集初始化成功。之后连接字符串写成mongodb://ip:27017,ip:27018/?replicaSetrs0这种格式驱动会自动发现主节点。写数据之前先建好索引。常用语句db.doc_event.createIndex({ userId: 1, createTime: -1 }) db.doc_event.createIndex({ expireAt: 1 }, { expireAfterSeconds: 0 })第一个是查询索引第二个是 TTL 索引MongoDB 会定期删除expireAt早于当前时间的文档。对于日志类文档存储TTL 索引是最省心的数据清理方案不用再写定时任务去删数据。3.3 配套的可视化与管理工具单靠命令行排查问题效率太低我一般会配几个可视化工具这里覆盖大家经常搜的几个。Kafka 侧常用的开源可视化工具包括AKHQ、Kafka UI和CMAK。AKHQ 的优势是不仅能看 topic、分区、消费组还能查看和提交 Connector 任务状态。它通过 Docker 启动非常快配置里指定 Kafka 地址即可。排查 Connector 故障时打开 AKHQ 找到对应 connect 集群直接看任务日志和状态比在服务器上翻日志方便得多。还有一个常用命令是bin/kafka-consumer-groups.sh --describe --group 消费组名 --bootstrap-server 地址可以看每个分区的LAG值。这个命令要养成习惯应用偶发延迟或积压时先看 LAG 是哪个分区高再决定加消费者还是查数据倾斜。MongoDB 侧Compass是官方图形客户端可视化查看文档、辅助构建索引、查看性能分析都方便。团队规模大、要监控多套环境可以使用Ops Manager它覆盖部署、监控、备份是比较完整的运维平台。日常快速看负载我常用mongostat和mongotop命令一个看整体读写指标一个看各集合的操作耗时。这些工具不复杂但在真正出问题时能省下大把时间。我见过不少团队连消费组的 LAG 都不看等 MongoDB 慢查询积了一大堆才去排查那时候链路已经堵到不可收拾了。4. 核心链路实操从 Producer 到 Mongo 落库4.1 Producer 端如何发消息决定了下游好不好写生产端是链路的起点消息带什么 key、按什么结构发直接影响下游消费和落库。我的经验是无论数据来源多乱进了 Kafka 的消息统一转成标准化 JSON里面保留业务唯一字段和原始数据字段方便后续解析和索引。一个 Java Producer 的典型配置Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka1:9092,kafka2:9092,kafka3:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.LINGER_MS_CONFIG, 20); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); props.put(ProducerConfig.ACKS_CONFIG, 1); KafkaProducerString, String producer new KafkaProducer(props); producer.send(new ProducerRecord(order-event, orderId, jsonString));重点解释下这几个参数linger.ms设成 20 表示发送端会攒 20 毫秒再打包发送配合batch.size能显著提高吞吐acks1表示只要 leader 写入成功就返回吞吐高但极端情况可能丢消息如果对可靠性要求很高用acksall配合min.insync.replicas2但写入延迟会高一些。文档存储场景如果允许极端情况下少量数据由 Mongo 端幂等兜底用acks1可以减少生产端拖后腿的概率。如果消息里的业务唯一键很明确比如orderId发送时务必带上 key。带上 key 有三个作用相同 key 进同一分区保证顺序落到 MongoDB 后可以直接用这个 key 做 _id 或 upsert 条件排查问题时能按 key 精确找消息不用全链路翻数据。4.2 首选方案Kafka Connect Mongo Sink Connector如果消息不需要复杂加工直接用MongoDB 官方提供的 Kafka Connect Sink Connector就能完成任务。它本质上是把 Kafka topic 里的每条消息解析成文档批量写入指定集合。这里给出一个可用配置{ name: mongo-doc-sink, config: { connector.class: com.mongodb.kafka.connect.MongoSinkConnector, tasks.max: 3, topics: order-event, connection.uri: mongodb://user:passmongo1:27017,mongo2:27018/?replicaSetrs0, database: bigdata, collection: order_doc, document.id.strategy: com.mongodb.kafka.connect.sink.processor.id.strategy.PartitionStrategy, key.converter: org.apache.kafka.connect.storage.StringConverter, value.converter: org.apache.kafka.connect.storage.StringConverter, value.converter.schemas.enable: false, writemodel.strategy: com.mongodb.kafka.connect.sink.writemodel.strategy.ReplaceOneDefaultStrategy, max.batch.size: 500, max.num.retries: 3 } }有几个配置要特别说明。tasks.max我设置成和 topic 分区数接近Connector 会拆出多个 task 并行消费document.id.strategy决定文档的 _id 怎么生成我建议不要用默认的自动 ObjectId而是用ProvidedStrategy直接取消息 key 作为 _id这样即使同一个订单被重复消费多次也能通过 ReplaceOne 策略去重更新writemodel.strategy用ReplaceOneDefaultStrategy配合幂等 _id效果等同于“有则更新无则插入”。启动 Connector 的方式有 REST API 和配置文件两种。生产环境用 REST API 更常见curl -X POST http://localhost:8083/connectors \ -H Content-Type: application/json \ -d mongo-sink-config.json启动后一定要做两件事查看 Connector 状态确认 task 都变RUNNING再往一个测试 topic 里发几条消息去 MongoDB 里查集合文档验证字段和 _id 是否符合预期。不要直接切线上 topic那个时期你会感谢自己多花了这十分钟。4.3 更灵活的备选自研消费者 bulkWriteConnector 虽然省事但遇到字段要清洗、多条消息要聚合、不同消息写不同集合它就不太行了。这时候我会写一个自研消费者核心逻辑很简单从 Kafka 拉一批消息解析成文档对象用 MongoDB 的批量接口一次写入。Java 侧的消费端和写入逻辑大致如下try (KafkaConsumerString, String consumer new KafkaConsumer(props)) { consumer.subscribe(Collections.singletonList(order-event)); while (running) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); ListDocument docs new ArrayList(); for (ConsumerRecordString, String record : records) { Document doc Document.parse(record.value()); doc.put(_id, record.key()); doc.put(eventTime, Instant.parse(doc.getString(eventTime))); docs.add(doc); if (docs.size() 500) { writeToMongo(docs); docs.clear(); } } if (!docs.isEmpty()) { writeToMongo(docs); } } } void writeToMongo(ListDocument docs) { ListWriteModelDocument writes docs.stream() .map(d - (WriteModelDocument) new ReplaceOneModel( Filters.eq(_id, d.get(_id)), d, new ReplaceOptions().upsert(true))) .collect(Collectors.toList()); collection.bulkWrite(writes, new BulkWriteOptions().ordered(false)); }这段代码有几个细节很关键。bulkWrite比逐条insertOne快很多批量越大 MongoDB 端落盘越省一般控制在 300 到 1000 条之间。用ReplaceOneModel构造 upsert以 _id 为匹配条件天然实现幂等重复消费同一批消息时后续只会覆盖文档不会产生重复记录。ordered(false)让批量里的单条失败不影响其他文档写入避免一条脏数据拖慢整批。批量大小不是越大越好。文档越大一条批量请求的 BSON 体积就越大超过 16MB 或默认的maxWriteBatchSize会被服务端拒绝。我的经验是单个文档 1KB 左右时500 条一批比较合适文档到 10KB 以上时降到 100 条一批更稳。这个值上线前最好压一下找到吞吐和失败率的平衡点。5. 常见问题与排查实录5.1 Kafka 侧连接、延迟与重复消费先聊连接问题。Producer 报连接超时大多数情况不是代码问题而是防火墙规则没放通9092 端口或者bootstrap.servers里写的是localhost而客户端从外部访问根本解析不到。另外Kafka 对外宣传的地址由advertised.listeners决定很多集群配置忘了改这个参数外部客户端连上了内网地址直接卡死。先去查这个配置再查防火墙。我在一台新部署的集群上踩过这个坑advertised.listeners没改成对外 IP结果外部所有客户端都连接超时日志里完全看不出来最后还是抓包发现的。消息延迟高是另一个高频场景。如果消费组 LAG 越来越高先看数据量是否突增再看单条消息处理是否变慢。我之前遇到一次消费延迟排查半天发现是某个业务方发了大批 5MB 的大 JSON单个文档解析、序列化、写入全链路耗时被拉高。解决办法是拆分消息或者限制最大消息大小Kafka 侧可以配置message.max.bytes来控制单条消息上限。重复消费的问题非常讨厌根因一般是消费者处理完消息、落库成功但在提交位移前进程挂了或者网络出问题重启后从旧位移重新拉取。应对思路就是那句老话消费端要做到幂等。用 key 做 _id、用 upsert 写入重复消费就是多执行一次覆盖写对业务无影响。这也从另一个角度说明了前面设计 _id 时选业务唯一键有多重要。5.2 MongoDB 侧连接池与写入并发MongoDB 端的经典故障是连接数耗尽。默认连接池上限是 100Java 驱动里是maxPoolSize100高并发消费时很容易打满然后出现获取连接超时消费线程集体阻塞LAG 直线上升。解决方法是根据并发线程数调大连接池例如消费线程 20 个连接池至少给到 50 到 100同时确保服务端net.maxIncomingConnections足够大。还有一个容易被忽略的参数是writeConcern。默认w1表示主节点写入成功即返回速度最快设置为wmajority表示要等大多数副本节点确认能避免主节点故障时丢数据但延迟高不少。文档存储场景如果允许在故障窗口内丢少数未确认的写入就用默认的w1如果每条数据都不能丢那就majority并做好写入延迟升高的准备。慢查询会让写入看起来像卡死。注意看操作日志里 BSON 大小和索引情况。我之前排查过一次“偶尔一条消息要写几秒”的问题最后发现是大量文档未覆盖索引更新或查询触发了全集合扫描。解决方法和早期规划文档模型时的建议一样上线前把查询和更新路径都列出来给每个路径配好索引真到积压时再补索引会非常痛苦。5.3 数据一致性去重与幂等落库数据一致性问题在分布式链路里很常见。Kafka 有“至少一次”的语义重复消费是常态不是小概率事件MongoDB 落库如果按每条数据 insert几乎必然出现重复文档。最彻底的方案是在写库时对业务唯一键做约束。MongoDB 虽然不像关系型数据库那样有传统主键约束但唯一索引 upsert能达到同样的效果。比如我们要求每个订单只有一条文档就建唯一索引{ orderId: 1 }写入时用ReplaceOneModel加upsertcollection.updateOne( Filters.eq(orderId, orderId), Updates.set(data, jsonValue), new UpdateOptions().upsert(true) );注意别把自动生成的_id当作业务主键_id是 MongoDB 给的业务方无法预知要用orderId这类业务唯一字段做条件。还有一类问题是消息里字段类型不一致。比如同一个字段有时是字符串 123有时是数字 123MongoDB 在 upsert 时可能报Cannot convert错误。经验做法是在消费者里做一次字段类型清洗统一数字字段为 long、时间字段为 Date别把类型问题留给 Mongo 侧。5.4 一条实战排查路径遇到“Kafka 有数据MongoDB 没数据”这类问题我总结了一条排查路径按顺序来基本能定位看消费组状态kafka-consumer-groups.sh --describe确认消费者是否在线LAG 是否在涨。看 Connector 状态如果是 Connector 模式用 AKHQ 或GET /connectors/{name}/status看 task 状态和错误信息。看消费者日志捕获异常类型比如反序列化失败、Mongo 写入超时、连接池耗尽。用测试消息复现向原 topic 发一条标准消息观察链路走到哪一步断了。查 MongoDB 慢日志和mongostat确认是写不进去还是写得太慢。这套流程基本能覆盖九成问题。最怕的是跳过前面步骤直接去翻 Mongo 日志容易绕弯路。排查分布式链路必须养成“从上游查到底游”的习惯。6. 性能调优与我的实操心得6.1 一份可直接抄走的吞吐参数清单把参数列成清单对照场景去匹配比记一堆理论值更实用。层面参数推荐值备注Kafka Producerlinger.ms10~50攒消息时间越大批量越大延迟也越高Kafka Producerbatch.size16~32KB单批最大字节数结合 linger 调整Kafka Consumermax.poll.records200~800一次 poll 拉取条数影响批量大小Kafka Consumerfetch.min.bytes1~5MB拉取最小字节提高拉取效率Mongo 写入bulkWrite批量条数300~1000视文档大小而定Mongo DrivermaxPoolSize50~200与消费线程数匹配不要盲目开大MongowriteConcern默认 w1 或 majority按可靠性要求取舍Mongo索引覆盖查询/更新/排序字段这是压测前必做项这些值不一定是绝对最优但作为起点够用。实际项目里我通常先按这个清单配置一遍再压测加监控日志根据真实吞吐做微调。压测时注意监控两个指标Kafka 消费侧的 LAG 和 MongoDB 侧的mongostat这两头任何一边出现瓶颈都要回到参数清单里去调对应项。6.2 消费延迟高与数据倾斜处理消费延迟高除了链路性能问题也可能出在数据倾斜上。按 key 分区很容易导致某些 key 的消息特别多比如一个热门商家占了 80% 的消息量对应分区的消费速度跟不上整体 LAG 下不去。排查数据倾斜可以先看每个分区的 LAG。如果某个分区 LAG 明显高于其他分区基本可以确认倾斜。解决办法有几个如果是 key 设计导致的热点可以考虑在 key 后面拼接随机后缀把数据打散到不同分区但前提是不要求同一个 key 消息保序如果 topic 分区数本身不够可以按新分区数重建 topic重新路由消息然后切换消费组还有一种方案是增加 Consumer 线程让更多线程消费同一个分区但这需要在代码里手动分配分区复杂度会上升。我曾经处理过一个订单事件倾斜的真实案例一个头部商户贡献了七成消息量导致对应分区消费者长期高负载。最后方案是把这个商户的 key 加后缀拆成四个 key订单顺序要求通过数据内的时间戳字段在 Mongo 侧排序来兜底整体 LAG 才降下来。所以做分区设计时一定要对头部 key 的占比有预判。6.3 监控与压测经验我之前讲到过“没有监控调优就是盲调”。推荐组合是Kafka 侧kafka-consumer-groups.sh配合 AKHQ/Kafka UI 看 LAG如果有 Confluent Control Center 当然更好没有就用这俩也够用。MongoDB 侧mongostat看读写速率mongotop看各集合耗时图形界面看 Compass大规模集群用 Ops Manager。链路整体搭一个简单的监控面板至少能看到“消息生产速率”和“Mongo 落库速率”两条曲线的对比。如果生产速率一直高于落库速率且没有收敛势头那就不是偶发延迟而是容量不足得扩充分区、消费者线程或者 Mongo 集群资源。压测时最容易被忽视的是文档大小分布。平均 1KB 和平均 10KB 两个场景同样一批消息的落库耗时差好几倍。压测要尽量模拟真实文档大小不要用统一小文档糊弄过去那样的压测结果没有任何参考意义。6.4 一些小而重要的设计教训这些教训放在最后是因为每一个都是用实际故障换来的。一是不要随便改 Kafka 的分区数。虽然 Kafka 支持扩大分区数但分区数量变更后哈希分配全部重算原来按 key 的顺序关系会被打乱。如果你的业务依赖分区内顺序扩大分区相当于破坏顺序保证。所以 Topic 分区数在设计阶段就要尽量贴近发展规划宁可比现实需求大一倍也不要因为懒而设 1 个分区。二是不要在消费者里做长耗时操作比如逐条调外部接口、做高复杂度计算。消费者线程卡住了poll 心跳就会超时触发 rebalance甚至导致二次消费和消息积压。长耗时操作要拆出去用独立线程池或者走 stream 处理链路别跟消息落库耦合在一起。三是存储层的写入是批量思维的不是请求思维的。单条 insert 的吞吐放到大数据场景完全不够看写入必须走批量接口。你在 Kafka 消费端攒一批数据再落库比逐条实时写要稳得多这个“攒”的过程就是给 MongoDB 减负的过程。四是多环境配置一定要自动化。消息格式、连接串、topic 名这些配置不同环境差异很大手工维护很容易出错。我强烈建议把环境相关配置放到统一配置中心或者用 CI/CD 生成不同环境的配置文件杜绝“生产环境用的是哪个 collection”这种低级错误。结尾这套 Kafka 加 MongoDB 的协作架构我前前后后维护了两年多最大的感受是它不是一个需要天天炫技的架构但胜在稳定、灵活、好扩展。真正的难点从来不在工具本身而是你有没有想清楚文档模型、分区策略、幂等写入这些细节。只要在这些地方提前把功课做足Kafka 转 Mongo 这条管道跑起来以后会非常省心。最后再分享一个我个人的习惯上线任何一条新管道之前一定先用测试 topic 和测试集合跑一遍全链路发几条构造好的“奇怪数据”字段缺失、类型不对、超长字符串看看消费端是不是会崩溃。这套习惯帮我拦下了无数次生产事故。毕竟再怎么优化的管道也怕“线上第一笔数据就是脏数据”这种开场。
返回列表