ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

数据库与消息队列通信实战:从outbox到CDC的可靠链路

数据库与消息队列通信实战:从outbox到CDC的可靠链路 简介一份聚焦数据库进程间通信的MQ技术方案文档面向需要在不能改动既有代码的前提下为数据库增删改操作触发外部任务的开发与运维人员系统讲解借助消息队列、命名管道、ZeroMQ和MySQL插件实现进程间通信与远程过程调用的思路。内容先剖析传统定时轮询的局限再结合短信邮件发送、图片处理、身份证号码校验、网站静态化、多库数据同步等场景给出具体触发方式并附有插件接口设计、触发器配置和同步异步调用示例。资源为单个Word文档压缩包大小一百一十六KB结构紧凑便于通读。已有九十七人学习浏览适合后台开发、数据库管理员参考。借助该方案可有效降低系统耦合提升数据变更响应的实时性为分布式系统提供灵活高效的信息传递机制。1. 数据库和业务进程之间为什么非要塞一个 MQ做过数据库同步、订单流转或者库存扣减的兄弟大概率都遇过这个场景业务逻辑写着写着发现“先写数据库还是先发 MQ”变成了一个哲学问题。先写库万一消息没发出去下游不干活先发 MQ消息出去了事务回滚下游拿着一条根本不存在的订单去处理。这个矛盾在热词里被反复搜成“先写数据库 先写mq”其实背后就是数据库进程间通信最典型的一个痛点——两个进程之间要协作但谁先谁后都会出现数据不一致。数据库进程间通信IPC本身不是新东西Unix 时代就有管道、信号量、共享内存、Socket 这些方案。但一旦进程里有一方是数据库或者消息要跨机器、跨语言、跨网络传递传统 IPC 就不够看了管道只能本机、共享内存要处理锁和生命周期、直接 Socket 长连接要自己管重试和积压。这时候把 MQMessage Queue插在中间本质上是把“谁先谁后”的竞争问题改造成“反正早晚都会到”的排队问题。MQ 能解决的场景很明确削峰填谷、解耦、异步化、失败重试。放到数据库场景里最常见的是三种——业务库把变更事件发出去给下游同步、多个应用进程通过 MQ 协调对同一张表的写入、以及把数据库里的慢查询或批量任务丢到队列里排队执行。这篇文章不聊概念直接用一套可复现的方案把 MQ 在数据库进程间通信里的选型、落地和排错讲透特别是那个让无数人翻车的“先写库还是先写 MQ”的时序问题。2. 先想清楚链路形态MQ 到底插在哪两个进程之间很多读者拿着 MQ 就去连数据库其实第一步错在没分清自己属于哪一种通信链路。链路形态不同表结构设计、消息载体、消费逻辑完全是三套写法。2.1 业务进程和数据库之间的异步写入链路第一种形态是业务进程比如一个订单服务要写数据库同时要通知另一个进程比如库存服务去扣减。常见做法是把“写库”和“发消息”放在同一个业务动作里但这两个动作跨越了进程边界无法共享同一个本地事务。我一般会把订单保存和消息发送拆成两个进程动作订单服务先把自己的业务落库然后发一条“订单已创建”的消息到 MQ库存服务作为消费者去消费。这里的核心问题是业务进程和 MQ 之间的可靠传递。进程可能写完库就崩了消息没发出去也可能消息发出去了但业务进程在事务提交前一瞬间崩了下游拿到的是脏数据。解决这个问题的关键不是疯狂重试而是通过消息表定时补偿来收敛。也就是说业务进程写业务表的同时往本地一张 outbox 表里写一条待发送的消息记录这两个写操作在同一个数据库事务里完成然后由另一个扫描进程把 outbox 里没发出去的消息捞出来发给 MQ。这是目前处理“先写库还是先发 MQ”最可靠的做法比本地消息表方案更省事也比事务消息的接入成本低得多。2.2 数据库主从或异构库之间的数据同步链路第二种更常见的形态是源数据库的变更要同步到另一个数据库异构数据库、数仓、缓存或者搜索引擎。常见做法是部署 Debezium、Canal 这类 CDC 工具抓 binlog把变更事件写入 MQ下游消费者拿到事件后回放到目标库。这种链路里 MQ 的作用不是解耦业务而是做速率适配和故障缓冲。源库的 TPS 可能有峰值目标库的写入能力可能跟不上如果没有 MQCDC 工具只能阻塞在源库上最终拖垮主库。而 MQ 把变更事件落盘排队目标库按自己的速度消费源库完全不感知下游压力。这个场景有一个关键设计消息里带的是“数据变更内容”而不是“业务指令”。比如用户表的一行 update消息体应该携带主键、变更前的镜像、变更后的镜像而不是“把小明年龄改成 25”这种业务化描述。因为下游可能是多份数据副本每份副本的过滤规则不同业务化描述只能满足一种消费逻辑而数据镜像谁都能用。2.3 选型判断原生 IPC 和 MQ 的边界在哪里实时性要求微秒级、只在本机两个进程间通信、数据量固定用共享内存或者 Unix Socket 就够了硬上 MQ 反而增加延迟和运维成本。但只要是“数据库里的数据要跨进程流动”这个前提MQ 几乎总是更合适的方案。选型上数据库同步场景我优先推 Kafka 或 RocketMQ因为消息量大、需要按 key 分区保序、需要长时间堆积。业务解耦场景用 RabbitMQ 或 RocketMQ 都行看团队熟悉程度。Riak、PostgreSQL 自带的一些 LISTEN/NOTIFY 机制适合轻量通知但不适合做可靠消息投递因为不是真正的落盘消息队列。3. 用 RabbitMQ 跑通数据库与业务进程间的 MQ 通信最小可复现方案选定 RabbitMQ 作为例子是因为它部署轻、路由灵活而且在“数据库 进程间通信”这个场景里最容易理解队列、交换机、路由键这些概念。下面这套方案解决的是 2.1 里的业务进程异步化链路我会把关键代码和参数都写清楚。3.1 数据库端的 outbox 消息表设计与写入先建一张 outbox 表注意这是一张业务库里的普通表不是 MQ 里的队列。这张表的职责是暂存“业务事件”让写业务表和写消息这件事在同一个数据库事务里原子完成。-- 建表业务库里的 outbox 消息表 CREATE TABLE outbox_message ( id BIGINT AUTO_INCREMENT PRIMARY KEY, aggregate_type VARCHAR(64) NOT NULL COMMENT 业务对象类型如 ORDER, aggregate_id VARCHAR(64) NOT NULL COMMENT 业务主键, event_type VARCHAR(64) NOT NULL COMMENT 事件类型如 ORDER_CREATED, payload JSON NOT NULL COMMENT 消息体JSON 格式, status TINYINT NOT NULL DEFAULT 0 COMMENT 0-待发送 1-已发送 2-发送失败待补偿, retry_count INT NOT NULL DEFAULT 0, next_retry_time DATETIME NOT NULL, create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, KEY idx_status_retry (status, next_retry_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;字段设计里的关键在于aggregate_type和aggregate_id。这两个字段决定了消息的去重维度——消费者拿到消息后用这个维度判断自己是否已经处理过这条数据。payload直接存 JSON也是为了方便下游直接透传不用反序列化两次。status 字段是给扫描程序用的0 表示待发送1 表示已发送但未确认2 表示多次失败进了死信。3.2 发送端事务内写表 事务后发消息的 Java 代码这里用一个 Spring Boot 风格的伪代码来演示核心逻辑。注意发消息的动作一定不能放在数据库事务里——网络抖动会把事务拖死而且事务还没提交消费者就收到消息此时读库会读不到。所以顺序必须是先提交事务再发消息。Transactional public void createOrder(OrderDTO dto) { // 1. 先写业务表和 outbox 消息在同一个事务里 orderMapper.insert(dto); outboxMapper.insert(OutboxMessage.builder() .aggregateType(ORDER) .aggregateId(dto.getOrderNo()) .eventType(ORDER_CREATED) .payload(JSON.toJSONString(dto)) .status(0) .nextRetryTime(LocalDateTime.now()) .build()); } EventListener(phase TransactionPhase.AFTER_COMMIT) public void afterCommit(OrderCreatedEvent event) { // 2. 事务提交成功后才发消息 OutboxMessage message outboxMapper.selectByAggregate(event.getOrderNo()); try { rabbitTemplate.convertAndSend(order.exchange, order.created, message.getPayload()); // 3. 标记为已发送此时消息已进入 RabbitMQ outboxMapper.markSent(message.getId()); } catch (Exception e) { // 发消息失败不要回滚业务事务留给补偿器处理 log.error(send mq failed, id{}, message.getId(), e); } }这段代码有两点需要重点解释。第一EventListener搭配TransactionPhase.AFTER_COMMIT是保证“事务提交后再发消息”的标准写法比在方法末尾手动发消息更安全因为 Spring 的事务代理会保证只有在真正提交成功后才会触发。第二这里发消息失败只是 catch 住不做任何重试是为了避免阻塞主流程。outbox 表里 status 还是 0交给后台补偿程序处理。3.3 补偿器把漏网之鱼捞回来如果不写补偿器前面两段代码做完大概能保证 95% 的消息送达率另外 5% 会丢在网络抖动、进程崩溃、MQ 短暂不可用这些场景里。补偿器的逻辑很简单定时扫 outbox 表里 status0 或 status1 且超过一定时间没收到确认的消息重新发给 MQ。// Quartz 或 XXL-JOB每 30 秒执行一次 public void compensate() { ListOutboxMessage pending outboxMapper.scanPending(LocalDateTime.now(), 100); for (OutboxMessage message : pending) { if (message.getRetryCount() 5) { outboxMapper.markDead(message.getId()); continue; // 超过重试次数进死信人工介入 } try { rabbitTemplate.convertAndSend(order.exchange, order.created, message.getPayload()); outboxMapper.markSent(message.getId()); } catch (Exception e) { outboxMapper.incrementRetry(message.getId(), LocalDateTime.now().plusSeconds(30)); } } }补偿器的重试时间用指数退避第一次失败后 30 秒重试第二次 1 分钟第三次 2 分钟最多 5 次。为什么不用固定间隔因为 MQ 不可用通常持续一段时间固定间隔只会给 MQ 雪上加霜。标记为已发送status1的消息也有可能在 MQ 端丢失比如 RabbitMQ 持久化到一半节点宕机。严谨的方案是开启 publisher confirm收到 broker 的 ack 后再标记已发送。生产环境我建议打开 publisher confirm代价是每条消息多一次网络往返但换来的是“发出去 broker 落盘”的强保证。3.4 消费者手动 ACK 和幂等是底线消费者端的重点不是怎么收消息而是收完消息之后怎么和数据库打交道。强烈建议关掉 RabbitMQ 的自动 ACK改成手动确认并且消费逻辑必须是幂等的。RabbitListener(queues order.queue) public void onMessage(OrderCreatedMessage msg, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws Exception { String orderNo msg.getAggregateId(); // 幂等检查用业务主键查消费记录表 if (consumeLogMapper.exists(ORDER, orderNo)) { channel.basicAck(tag, false); return; // 已经处理过直接确认 } try { // 处理消息写自己的业务表或调用内部服务 stockService.deduct(msg.getPayload()); // 业务处理成功后记消费日志 确认消息这两个动作通常在同一个本地事务 consumeLogMapper.insert(ORDER, orderNo); channel.basicAck(tag, false); } catch (Exception e) { // 处理失败不确认也不重新入队等人工排查 channel.basicNack(tag, false, false); log.error(consume message failed, orderNo{}, orderNo, e); } }这里的核心参数在basicNack的第三个参数requeuefalse。很多新手在这里填 true结果消费失败的消息无限循环把日志刷爆。requeuefalse 的好处是消息会被 RabbitMQ 丢弃或进死信队列至少要保证日志里有完整记录而不是被循环消费淹没。幂等检查用独立的消费记录表不要依赖业务表本身做判断因为业务表可能因为各种原因被手工修改过消费记录表更纯粹。4. 数据库变更实时同步到 MQ用 Debezium 加 Kafka 搭一条完整链路如果你不是要“业务主动发消息”而是希望“数据库里的任何变化都被自动感知并同步出去”那就需要回到 2.2 的链路——用 CDC 工具抓 binlog 喂给 MQ。这是目前数据库同步软件和异构库同步的主流架构和 3.x 章节的方案是互补关系业务明确感知的用 outbox 主动发无法侵入业务代码的用 CDC 自动抓。4.1 架构选型为什么用 Debezium Kafka 而不用直连轮询直连轮询的思路是写一个定时任务每秒钟 SELECT 一次业务表把新增和修改的数据捞出来发到 MQ。这个方案在小数据量下能用但它有三个硬伤第一无法感知删除事件只能靠逻辑删除字段补偿第二SELECT 是“查一次快照”而不是“流式变更”每次都要全表或按更新时间扫描做不到秒级实时第三会对业务库产生额外查询压力大表场景直接拖垮数据库。Debezium 的思路完全不同它伪装成一个 MySQL 从库订阅 binlog 流数据库提交的任何变更insert、update、delete都会以事件形式流式推给 Debezium再由 Debezium 写入 Kafka。这个机制下源库零侵入也不会有轮询延迟。4.2 Docker 部署 Debezium 与 Kafka 的最小命令整套环境用 Docker Compose 最省事这里给出最小可用的部署文件version: 3 services: zookeeper: image: confluentinc/cp-zookeeper:7.4.0 environment: ZOOKEEPER_CLIENT_PORT: 2181 kafka: image: confluentinc/cp-kafka:7.4.0 ports: - 9092:9092 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_AUTO_CREATE_TOPICS_ENABLE: true connect: image: debezium/connect:2.4 ports: - 8083:8083 environment: BOOTSTRAP_SERVERS: kafka:9092 GROUP_ID: mq-connect-group CONFIG_STORAGE_TOPIC: my_connect_configs OFFSET_STORAGE_TOPIC: my_connect_offsets STATUS_STORAGE_TOPIC: my_connect_statuses部署完要等 Kafka 和 Debezium Connect 的日志不再报错然后用 REST API 注册一个 MySQL connector——这是启动数据抓取的关键步骤注册命令如下curl -X POST -H Content-Type: application/json --data { name: mysql-orders-connector, config: { connector.class: io.debezium.connector.mysql.MySqlConnector, database.hostname: 172.16.1.10, database.port: 3306, database.user: debezium, database.password: yourpassword, database.server.name: orders-server, database.include.list: app_db, table.include.list: app_db.t_order, database.history.kafka.bootstrap.servers: kafka:9092, database.history.kafka.topic: schema-changes.orders, topic.prefix: orders-db, key.converter: org.apache.kafka.connect.json.JsonConverter, value.converter: org.apache.kafka.connect.json.JsonConverter, value.converter.schemas.enable: false } } http://localhost:8083/connectors参数说明里最重要的一对是table.include.list和database.include.list。前者精确到表名后者精确到库名必须都配置否则 Debezium 会抓所有库所有表产生大量无用事件占用 Kafka 磁盘。database.server.name是逻辑名称它会成为 Kafka topic 名称的前缀例如这里会生成orders-server.app_db.t_order这个 topic。database.history.kafka.topic是 Debezium 用来记录所有 DDL 变更的专属 topic千万不能删除。删除之后 Debezium 丢失了表结构历史重启会直接起不来报错信息会让人一头雾水——这也是一个新手的常见坑。4.3 topic 与分区键的设计规则Debezium 默认按主键决定消息发往哪个分区也就是说同一个主键的变更事件会进入同一个分区保证了这个主键维度上的消息顺序。这个特性在同步场景里至关重要因为 binlog 里同一行的 update 是有先后顺序的如果散到不同分区消费端可能先处理后一次更新、再处理前一次目标库最终数据就是错的。如果一个业务表没有主键强烈建议在 MySQL 里补上。没有主键的表Debezium 无法生成稳定分区键顺序性完全丧失这对数据库同步是致命的。多表关联需要合并的场合可以在 Kafka 消费端再按业务 key 做一次本地聚合不要尝试在 Debezium 端改分区策略。消费方要记住Kafka只保证分区内有序不保证跨分区有序。如果一个业务表的变更量不足以分散到多分区就用单分区 topic 换顺序性性能完全够。4.4 消费端落库的回放策略与类型映射消费端拿到变更事件后写目标库最常翻车的环节是类型映射。Debezium 默认把 MySQL 的 datetime 转成微秒级时间戳把 decimal 转成带 scale 的字符串。如果你直接把这个 JSON 塞给目标库很可能出现日期差 8 小时、金额精度丢失的情况。我推荐的消费端策略是先把 JSON 转成内部的泛型 Map不做强类型绑定datetime 字段统一用yyyy-MM-dd HH:mm:ss格式化后再落库decimal 字段用字符串接收由目标库字段类型决定精度空的 update 事件数据没实际变化直接跳过不产生落库操作落库的顺序性方面消费端要严格控制并发度。单个 topic 分区的消费并发度应该等于 1也就是一个分区一个线程不要用线程池并发消费同一个分区的事件。否则即使 Kafka 端有序并发落库也会乱序。5. 数据库与 MQ 联调避坑从丢消息到死锁的 5 条血泪经验从 outbox 到 CDC整个链路跑通不难跑稳很难。下面这几条是我在项目和同行交流里反复遇到的坑每一条都对应一次真实翻车经历。全部按“现象 → 原因 → 解决”来写值得直接抄进团队的代码评审规范。5.1 事务消息与本地事务的死锁先发消息导致数据库连接池耗尽现象把发 MQ 的操作写在数据库事务内高峰期出现大量 database connection is not available 报错数据库连接数被打满业务接口大面积超时。原因每个事务持有数据库连接的同时还在等 RabbitMQ 的 broker 确认。MQ 一旦开始积压或网络抖动事务迟迟不释放连接。连接池默认 10~50 个连接几百个请求就把池子占完了。解决发消息永远在事务提交之后。用 Spring 的TransactionSynchronizationManager.registerSynchronization注册 afterCommit 回调或者在事务方法里手动提交事务再发消息。如果团队没有 Spring就用 3.2 里的 outbox 模式补偿器兜底。5.2 消费失败无限重试把日志刷爆requeue 参数用错现象某天凌晨 MQ 消费者持续打印同一批异常堆栈日志文件 1 小时涨了几个 GBKafka 的消费组 lag 一直不降。原因消费失败后调用了basicNack(tag, false, true)第三个参数requeuetrue让消息重新回到队列头部消费者立刻再次收到它形成死循环。如果异常是持久的比如 JSON 格式错误就会永远循环下去。解决不能用 requeue 来解决消费失败。生产环境统一用requeuefalse配合死信交换机把反复失败的消息隔离到专门队列由人工或定时任务处理。RabbitMQ 的x-dead-letter-exchange参数配置好之后失败消息自动转移。5.3 先写库后发消息事务提交瞬间宕机导致消息缺失现象发送端已经写完数据库并且事务成功提交但进程在调用 MQ 客户端之前崩溃业务数据入库了消息没发出去下游永远不知道这笔订单存在。原因数据库事务与 MQ 消息投递是两个没有共同原子性的动作。即使消息在事务后发送也存在“事务已提交、代码还没执行到发送”的窗口期。代码层面无论如何缩短这个窗口都不能消除它。解决outbox 表方案是当前最可靠的兜底。写入业务表的同时写 outbox 表两条 SQL 在同一个本地事务里即使进程崩溃outbox 表里仍有遗留数据补偿器扫描后会补发。不要相信“事务后立即发消息”能到 100%所有不依赖消息表的高可用方案都有流失窗口。5.4 数据库连接池大小与消费并发度不匹配消费者线程过多拖垮数据库现象消费者设置了 32 个并发线程每条消息都查一次数据库。数据库 CPU 飙升但其实每秒只处理了几百条消息大部分线程都在等连接。数据库连接池默认 20 个连接32 个线程抢 20 个连接一半线程在阻塞。原因MQ 的消费并发度不等于数据库的最大承载能力。很多团队把多线程当作吞吐的万能药忽略了数据库连接池、事务锁、SQL 执行耗时共同组成的真实瓶颈。解决消费者的并发线程数不应该超过数据库连接池大小的 50%。比如连接池 20消费者线程就只能开 10。如果目标是要提高消费吞吐先优化单条消息的 SQL 耗时再少量提升并发。同时要观察数据库的活跃连接数和线程等待时间这两个指标比消费者线程数更能说明问题。5.5 Kafka 消费端直接按表回放DDL 变更后事件结构与表结构不匹配现象数据库给某张表加了一个字段之后下游消费端一直在抛 unknown column 异常从 Debezium 发出的消息已经包含新字段但目标库表里没有。原因Debezium 的 schema history 里有历史表结构新事件带上了新增列但目标表的 DDL 没有同步或者同步了但消费端的 INSERT 语句还是硬编码旧字段列表。这属于典型的“源头变了、管道变了、终点没变”的不一致问题。解决目标库结构变更必须走自动化迁移工具手动改库不可靠。消费端落库不要写死字段列表用 JSON 里的字段动态构建 INSERT 语句未知字段先记录到扩展表不要因为多字段而失败。另外给 Debezium 开启 schema change event 的监听DDL 变更时自动触发目标库的对比迁。6. 积压怎么救、顺序怎么验证一套能直接落地的 MQ 可观测性配置最后这章写给要上生产的团队消息中间件装上只是第一步线上跑起来以后真正考验人的是积压排查和顺序性验证。我不讲大而全的监控平台搭建只给一套轻量但关键时刻能救命的排查套路。6.1 积压的三种判断方法与隔离方案先定义“积压”消息从生产者发出到被消费者拉取耗时超过正常基线的 10 倍。判断积压不能只看监控大盘要看三个具体指标消费组的 lag 值Kafka或队列堆积数RabbitMQ这个数字是积压的直接表达消费者的平均消费耗时如果单条消息耗时从 5ms 涨到 500ms即使堆积数没涨也要注意消息在 MQ 侧的停留时间这个指标比堆积数更准确因为堆积数可能因消息过期被清除积压发生后的第一件事不是扩容消费者而是确认瓶颈在消费者还是在下游数据库。用 5.4 的方法先看数据库活跃连接和慢 SQL。曾经有一回 Kafka lag 堆到几百万最后定位到是消费端的 SQL 少写了索引每次处理都全表扫描30 个消费者线程全部卡在数据库上再扩一倍也没用。如果确认瓶颈在消费者本身常规做法是临时增加消费组实例数。Kafka 的消费者组会自动做分区重分配RabbitMQ 则需要在 queue 的消费者线程数上做调整。两种操作都不需要改代码但要记住扩容只对 CPU 密集或 IO 等待型的消费逻辑有效如果瓶颈在数据库锁扩容只会增加锁竞争恶化问题。6.2 消息顺序性的验证脚本顺序性出问题是 MQ 场景里最隐蔽的故障表面上数据都对但对账时总差几条。最直接的验证方式是在消费者处理完消息后把消息里的序号写入日志或数据库然后用一段 SQL 脚本检查乱序。-- 验证同一业务主键的消息顺序是否错乱 SELECT aggregate_id, event_sequence, LAG(event_sequence) OVER (PARTITION BY aggregate_id ORDER BY create_time) AS prev_seq FROM consume_record WHERE create_time DATE_SUB(NOW(), INTERVAL 1 HOUR) HAVING prev_seq event_sequence;这段 SQL 的精髓在LAG窗口函数按消息创建时间排序后取上一条事件的序号如果上一条序号比当前这条大说明处理顺序反了。这个脚本建议做成定时任务每小时跑一次有异常就告警。拉出乱序记录后最常见的修复手段有两种如果乱序范围很小且不会影响最终一致性可以直接补一条修正消息如果影响较大需要把对应业务 key 的消费位点回退从乱序位置重新消费。回退位点是一件危险操作要先停止消费组再重置 offset确认无误后才能重启。6.3 压测时最容易骗人的参数prefetch 与 max.poll.records给 MQ 链路做压测时很多人上来就调大prefetchRabbitMQ或max.poll.recordsKafka其实这两个参数被严重高估了。prefetch 表示消费者预取多少条消息到本地缓存它只能提高网络吞吐但不能提高数据库处理能力。如果消费者的下一跳是数据库prefetch 调大只会让本地堆积更多待处理消息数据库跟不上时内存先爆。我常用的取值是消费者单条消息数据库耗时小于 10ms 时prefetch 设 50耗时 10~100ms 时prefetch 设 10~20耗时大于 100ms比如写数仓大批量导入prefetch 设 1~3。这个取值经验在 Kafka 侧对应max.poll.records同样遵循“消息处理耗时越长单次拉取条数越少”的原则。顺序验证、积压排查、连接池联动、参数取值——这些才是 MQ 在数据库通信场景里真正有门槛的部分。框架代码大家都写得出来区分方案是否可靠的是这些看不见的边界。我自己也经历过先写库后发消息的崩盘翻车从那以后 outbox 表和补偿器成了我所有 MQ 方案的标配。希望这些参数和避坑点能帮你少踩一次坑让这套链路一次跑稳。本文还有配套的精品资源点击获取
返回列表