削峰填谷与状态共享:华为云 DMS Kafka + DCS Redis 在 SpringBoot 中的实战

举报
行者·全栈架构师 发表于 2026/08/25 11:59:44 2026/08/25
【摘要】 本文复盘订单中心接入华为云 DMS for Kafka 与 DCS for Redis 的全过程。从异步解耦、削峰填谷、分布式锁、缓存一致性四个场景切入,覆盖消息可靠投递、幂等消费、死信队列、缓存穿透/击穿/雪崩防护、Redisson 分布式锁、Binlog 订阅缓存刷新等 7 个核心实现。包含 6 段可运行代码、4 张架构/流程图、5 个真实踩坑案例和压测前后对比数据

摘要: 本文复盘订单中心接入华为云 DMS for Kafka 与 DCS for Redis 的全过程。从异步解耦、削峰填谷、分布式锁、缓存一致性四个场景切入,覆盖消息可靠投递、幂等消费、死信队列、缓存穿透/击穿/雪崩防护、Redisson 分布式锁、Binlog 订阅缓存刷新等 7 个核心实现。包含 6 段可运行代码、4 张架构/流程图、5 个真实踩坑案例和压测前后对比数据(订单创建 P99 从 420ms 降到 86ms,库存扣减并发从 200 TPS 提升到 3500 TPS),并附 openGauss Binlog + 鲲鹏 NUMA 调优方案。

互动:你在 SpringBoot 项目里是怎么处理"下单高峰 + 库存扣减"这对矛盾的?用过哪些消息队列和缓存方案,踩过哪些坑?欢迎评论区交流。

一、为什么订单中心需要 Kafka + Redis?

1.1 同步架构的天花板

订单中心最早是纯同步架构:用户下单 → 扣库存 → 创订单 → 扣优惠券 → 发短信 → 加积分。一次下单串行 6 个 RPC,P99 耗时 420ms。大促时并发一上来,三个问题同时爆雷:

010-dms-dcs-middleware-practice_diagram_1.png

最严重的一次:短信供应商限流,导致下单接口大量超时,用户重试又加剧雪崩,30 分钟内下单成功率从 99% 跌到 62%。运营在群里发的截图我到现在都记得——满屏的"系统繁忙"。

1.2 异步 + 缓存的解药

引入 Kafka + Redis 后,下单链路重构为:

010-dms-dcs-middleware-practice_diagram_2.png

核心收益:用户感知的"下单"只剩 3 步(预扣 + 创单 + 发消息),P99 从 420ms 降到 86ms;后续 4 个动作异步消费,任一挂了不影响下单主链路。

1.3 为什么选华为云 DMS + DCS

选型考量 自建 Kafka/Redis 华为云 DMS/DCS
运维 自己管集群、扩容、故障切换 托管服务,免运维
高可用 需自己搭主从 + 哨兵 DCS 主备/集群自动切换,DMS 多副本
监控 自建 Prometheus + Exporter 内置监控看板,与 AOM 打通
安全 自配 SASL/SSL VPC 内网 + IAM 鉴权 + TLS
信创 DCS 兼容 openEuler,DMS 支持鲲鹏节点
成本 3 节点 ECS + 运维人力 按需付费,小规格月费 ~200 元起

关键认知:托管中间件最大的价值不是省机器钱,是省运维人力故障自动恢复。自建 Redis 哨兵切换我踩过 3 次脑裂,DCS 主备自动切换上线后再没出过。

二、华为云 DMS + DCS 整体架构

2.1 订单中心中间件架构

010-dms-dcs-middleware-practice_diagram_3.png

2.2 资源规格与成本

资源 规格 用途 月费估算
DMS Kafka kafka.2u4g.cluster × 3 order/stock/topic ~600 元
DCS Redis redis.ha.xu2.t1.r2.2GB 分布式锁 + 缓存 ~280 元
DCS Redis redis.cluster.xu2.t1.r2.4GB × 6 库存预扣(高频写) ~720 元
Canal ECS 2C4G Binlog 订阅 ~120 元

三、DMS Kafka 实战:可靠投递与幂等消费

3.1 SpringBoot 接入华为云 DMS

DMS Kafka 兼容开源 Kafka 协议,接入只需替换 bootstrap-servers 和 SASL 配置:

spring:
  kafka:
    bootstrap-servers: 192.168.0.10:9096,192.168.0.11:9096,192.168.0.12:9096
    properties:
      security.protocol: SASL_SSL
      sasl.mechanism: PLAIN
      sasl.jaas.config: >-
        org.apache.kafka.common.security.plain.PlainLoginModule required
        username="DMS_USER" password="DMS_PWD";
      ssl.truststore.location: classpath:dms-truststore.jks
      ssl.truststore.password: changeit
    producer:
      acks: all
      retries: 3
      max.in.flight.requests.per.connection: 1
      enable.idempotence: true
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
    consumer:
      group-id: order-cg
      enable-auto-commit: false
      max-poll-records: 100
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer

关键参数说明

  • acks=all + enable.idempotence=true:生产者精确一次语义,防消息丢失 + 防重复
  • max.in.flight.requests.per.connection=1:保证重试时不乱序
  • enable-auto-commit=false:手动提交 offset,消费成功才提交,防消息丢失

3.2 可靠投递:本地事务 + 消息表

下单发消息最大的坑是"DB 提交了但消息没发出去"。我用"本地事务 + 消息表"方案:

@Service
@RequiredArgsConstructor
public class OrderService {

    private final OrderRepository orderRepository;
    private final OutboxMessageRepository outboxRepository;

    @Transactional
    public Long create(OrderCreateDTO dto) {
        Order order = buildOrder(dto);
        orderRepository.insert(order);

        OutboxMessage msg = OutboxMessage.builder()
                .aggregateId(order.getId())
                .topic("order.created")
                .payload(JSON.toJSONString(order))
                .status(OutboxStatus.PENDING)
                .createdAt(TimeUtil.now())
                .build();
        outboxRepository.insert(msg);

        return order.getId();
    }
}

@Component
@RequiredArgsConstructor
public class OutboxPublisher {

    private final OutboxMessageRepository outboxRepository;
    private final KafkaTemplate<String, String> kafkaTemplate;

    @Scheduled(fixedDelay = 200)
    public void publish() {
        List<OutboxMessage> pending = outboxRepository.findTop100ByStatus(OutboxStatus.PENDING);
        for (OutboxMessage msg : pending) {
            try {
                kafkaTemplate.send(msg.getTopic(), msg.getAggregateId().toString(), msg.getPayload()).get();
                outboxRepository.updateStatus(msg.getId(), OutboxStatus.SENT);
            } catch (Exception e) {
                log.warn("消息发送失败,下次重试: {}", msg.getId(), e);
            }
        }
    }
}

原理:订单和消息表在同一个事务里写入,要么都成功要么都失败;定时任务扫消息表投递到 Kafka,发送成功才标记 SENT。这就是 Outbox 模式,比"先发消息再写库"或"先写库再发消息"都可靠。

3.3 幂等消费:防重复处理

Kafka 至少投递一次,消费端必须幂等。订单创建消息的消费逻辑:

@Component
@RequiredArgsConstructor
public class OrderCreatedConsumer {

    private final StockService stockService;
    private final IdempotentRepository idempotentRepository;

    @KafkaListener(topics = "order.created", groupId = "order-cg")
    public void onOrderCreated(ConsumerRecord<String, String> record, Acknowledgment ack) {
        OrderEvent event = JSON.parseObject(record.value(), OrderEvent.class);
        String idempotentKey = "order:consume:" + event.getOrderId();

        try {
            boolean firstTime = idempotentRepository.insertIfAbsent(idempotentKey, record.offset());
            if (!firstTime) {
                log.info("重复消息,跳过: orderId={}, offset={}", event.getOrderId(), record.offset());
                ack.acknowledge();
                return;
            }

            stockService.deduct(event.getOrderId(), event.getItems());

            ack.acknowledge();
        } catch (Exception e) {
            idempotentRepository.delete(idempotentKey);
            log.error("消费失败,不提交 offset 等重试: orderId={}", event.getOrderId(), e);
        }
    }
}

幂等表 SQL(基于 openGauss):

CREATE TABLE idempotent_record (
    idempotent_key VARCHAR(128) PRIMARY KEY,
    kafka_offset   BIGINT NOT NULL,
    created_at     TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

CREATE INDEX idx_idempotent_created ON idempotent_record(created_at);

踩坑①:幂等表无限增长
早期没清理幂等表,3 个月涨到 2 亿行,insert 变慢。后来加定时任务清理 7 天前的记录,且 key 用 业务前缀:订单ID(订单 ID 本身唯一),表大小稳定在 50 万行。

3.4 死信队列:消费失败的兜底

消费重试 3 次仍失败的消息投递到死信 topic,人工介入:

@Component
@RequiredArgsConstructor
public class DeadLetterHandler {

    private final KafkaTemplate<String, String> kafkaTemplate;
    private final RetryCounterRepository retryCounterRepository;

    public void handle(ConsumerRecord<String, String> record, Exception e) {
        int retryCount = retryCounterRepository.increment(record.key());
        if (retryCount >= 3) {
            DeadLetterMessage dlm = DeadLetterMessage.builder()
                    .originalTopic(record.topic())
                    .originalPartition(record.partition())
                    .originalOffset(record.offset())
                    .key(record.key())
                    .value(record.value())
                    .errorMessage(e.getMessage())
                    .build();
            kafkaTemplate.send("order.dead-letter", record.key(), JSON.toJSONString(dlm));
            log.error("消息进入死信队列: key={}, retry={}", record.key(), retryCount);
        } else {
            throw new RuntimeException("等待 Kafka 重试", e);
        }
    }
}

死信 topic 接入企业微信告警,运维同学收到告警后到 DMS 控制台查看消息内容,决定是修代码重投还是丢弃。

四、DCS Redis 实战:库存预扣与分布式锁

4.1 库存预扣:Lua 脚本保证原子性

高并发下单最大的敌人是超卖。用 Redis Hash 存库存,Lua 脚本原子预扣:

-- stock_deduct.lua
-- KEYS[1]: 库存 hash key
-- ARGV[1]: SKU ID
-- ARGV[2]: 扣减数量
-- ARGV[3]: 预扣订单号(用于回滚)

local stockKey = KEYS[1]
local sku = ARGV[1]
local qty = tonumber(ARGV[2])
local orderId = ARGV[3]

local current = tonumber(redis.call('HGET', stockKey, sku))
if current == nil then
    return -1
end
if current < qty then
    return -2
end

redis.call('HINCRBY', stockKey, sku, -qty)
redis.call('HSET', stockKey .. ':prehold', orderId .. ':' .. sku, qty)
return current - qty

SpringBoot 调用:

@Service
@RequiredArgsConstructor
public class StockPreDeductService {

    private final StringRedisTemplate redisTemplate;
    private DefaultRedisScript<Long> deductScript;

    @PostConstruct
    public void init() {
        deductScript = new DefaultRedisScript<>();
        deductScript.setScriptSource(new ResourceScriptSource(
                new ClassPathResource("lua/stock_deduct.lua")));
        deductScript.setResultType(Long.class);
    }

    public DeductResult preDeduct(String sku, int qty, String orderId) {
        Long remain = redisTemplate.execute(
                deductScript,
                Collections.singletonList("stock:warehouse:001"),
                sku, String.valueOf(qty), orderId);

        if (remain == null || remain == -1L) {
            return DeductResult.notFound();
        }
        if (remain == -2L) {
            return DeductResult.insufficient();
        }
        return DeductResult.success(remain);
    }
}

压测对比(4C8G 单实例,wrk 200 并发):

方案 TPS P99 超卖率
DB 直接 UPDATE stock SET qty=qty-? WHERE sku=? AND qty>=? 220 380ms 0%
DB + 行锁 SELECT ... FOR UPDATE 180 520ms 0%
Redis Lua 预扣 3,500 18ms 0%
Redis 无 Lua(先读后写) 4,100 15ms 1.7%

关键数据:Redis Lua 预扣比 DB 直接扣减快 16 倍,且零超卖。无 Lua 方案虽然 TPS 更高但超卖率 1.7%,不可接受。

4.2 分布式锁:Redisson 替代手写 SETNX

库存最终扣减(DB 层)需要分布式锁防并发。早期手写 SETNX + EXPIRE 踩过坑:

// 错误写法:SETNX 和 EXPIRE 不是原子操作
Boolean locked = redisTemplate.opsForValue().setIfAbsent(key, "1");
if (locked) {
    redisTemplate.expire(key, 10, TimeUnit.SECONDS);  // 若进程在这里挂了,锁永不释放
}

改用 Redisson:

@Configuration
public class RedissonConfig {

    @Bean(destroyMethod = "shutdown")
    public RedissonClient redissonClient() {
        Config config = new Config();
        config.useClusterServers()
                .addNodeAddress(
                        "redis://192.168.0.20:6379",
                        "redis://192.168.0.21:6379",
                        "redis://192.168.0.22:6379")
                .setPassword(System.getenv("DCS_PWD"))
                .setScanInterval(2000)
                .setIdleConnectionTimeout(10000)
                .setConnectTimeout(5000)
                .setRetryAttempts(3);
        return Redisson.create(config);
    }
}

@Service
@RequiredArgsConstructor
public class StockFinalDeductService {

    private final RedissonClient redissonClient;
    private final StockRepository stockRepository;

    public void finalDeduct(String sku, int qty, String orderId) {
        RLock lock = redissonClient.getLock("lock:stock:" + sku);
        try {
            if (!lock.tryLock(3, 30, TimeUnit.SECONDS)) {
                throw new BusinessException("库存扣减锁竞争失败: " + sku);
            }
            stockRepository.deduct(sku, qty, orderId);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new BusinessException("加锁中断", e);
        } finally {
            if (lock.isHeldByCurrentThread()) {
                lock.unlock();
            }
        }
    }
}

Redisson 的看门狗(watchdog)会自动续期,避免业务执行时间超过锁过期时间导致锁提前释放。

踩坑②:Redisson 看门狗在容器里不生效
早期用 lock.lock(30, SECONDS) 显式指定过期时间,看门狗不启动,结果一次 DB 慢查询导致锁过期被别人抢走,出现重复扣减。改成 lock.tryLock(waitTime, -1, SECONDS)(leaseTime=-1 启动看门狗自动续期)后解决。

4.3 缓存三大敌人:穿透、击穿、雪崩

订单详情查询的缓存防护:

@Service
@RequiredArgsConstructor
public class OrderQueryService {

    private final StringRedisTemplate redisTemplate;
    private final OrderRepository orderRepository;
    private final BloomFilter<String> orderBloomFilter;

    @Cacheable(value = "order:detail", key = "#orderNo", unless = "#result == null")
    public OrderDetailVO query(String orderNo) {

        if (!orderBloomFilter.mightContain(orderNo)) {
            return null;
        }

        String nullKey = "order:null:" + orderNo;
        if (Boolean.TRUE.equals(redisTemplate.hasKey(nullKey))) {
            return null;
        }

        Order order = orderRepository.findByOrderNo(orderNo);
        if (order == null) {
            redisTemplate.opsForValue().set(nullKey, "1", 60, TimeUnit.SECONDS);
            return null;
        }

        return toVO(order);
    }
}

三层防护

问题 触发场景 防护手段
缓存穿透 查询不存在的订单号 布隆过滤器拦截 + 空值缓存(60s)
缓存击穿 热点 key 过期瞬间大量请求 @Cacheable + 单飞机制(Redisson 互斥锁重建)
缓存雪崩 大量 key 同时过期 过期时间加随机扰动 ttl = base + random(0, 300s)

踩坑③:布隆过滤器容量预估错误
早期按订单总量 1000 万初始化布隆过滤器,结果第二年订单量到 1500 万,误判率从 0.1% 飙到 8%,大量不存在的订单号穿透到 DB。布隆过滤器容量要按 3 年预估,且误判率设 0.01% 才稳。

五、缓存一致性:Canal + Binlog 异步刷新

5.1 为什么不用先更新 DB 再删缓存

"先更 DB 再删缓存"在并发读写时仍会出现不一致:

010-dms-dcs-middleware-practice_diagram_4.png

更稳的方案是订阅 Binlog 异步刷新缓存:

010-dms-dcs-middleware-practice_diagram_5.png

5.2 Canal Client 实现

@Component
@RequiredArgsConstructor
public class CanalClientRunner implements CommandLineRunner {

    private final StringRedisTemplate redisTemplate;
    private final CanalConnector connector;

    @Override
    public void run(String... args) {
        connector.connect();
        while (true) {
            Message message = connector.getWithoutAck(1000);
            long batchId = message.getId();
            if (batchId != -1) {
                try {
                    process(message.getEntries());
                    connector.ack(batchId);
                } catch (Exception e) {
                    log.error("Canal 处理失败,回滚 batch", e);
                    connector.rollback(batchId);
                }
            }
        }
    }

    private void process(List<CanalEntry.Entry> entries) {
        for (CanalEntry.Entry entry : entries) {
            if (entry.getEntryType() != CanalEntry.EntryType.ROWDATA) continue;
            CanalEntry.RowChange change = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
            String table = entry.getHeader().getTableName();

            for (CanalEntry.RowData rowData : change.getRowDatasList()) {
                Map<String, String> after = parseColumns(rowData.getAfterColumnsList());
                String orderNo = after.get("order_no");
                if (orderNo != null) {
                    redisTemplate.delete("order:detail::" + orderNo);
                    log.debug("Canal 刷新缓存: table={}, orderNo={}", table, orderNo);
                }
            }
        }
    }
}

踩坑④:Canal 断连重连丢数据
Canal Server 重启时,Client 会从断点继续,但如果 Client 处理失败 rollback 后重试一直失败,会卡住整个 binlog 消费。我加了"重试 5 次仍失败则告警 + 跳过"策略,并把失败的 binlog 位点记录到 DB 表供人工补刷。

踩坑⑤:Canal 与 openGauss 兼容性
Canal 原生支持 MySQL Binlog,openGauss 兼容 PostgreSQL 用的是 WAL 不是 Binlog。最终用华为云 DRS(数据复制服务)替代 Canal,DRS 原生支持 openGauss → Redis 的数据订阅,配置更简单且由华为云托管。

六、鲲鹏 + 国产化适配

6.1 DCS Redis 在鲲鹏上的优化

DCS Redis 集群部署在鲲鹏 920 节点,开启 NUMA 绑核提升内存访问性能:

# 鲲鹏节点上 Redis 进程绑核
numactl --cpunodebind=0 --membind=0 redis-server /etc/redis/redis.conf

# 关键内核参数(openEuler)
sysctl -w net.core.somaxconn=65535
sysctl -w vm.overcommit_memory=1
echo never > /sys/kernel/mm/transparent_hugepage/enabled

6.2 性能对比

指标 x86 DCS 4GB 鲲鹏 920 DCS 4GB 提升
SET QPS 89,000 118,000 +32.6%
GET QPS 112,000 145,000 +29.5%
Lua 执行 QPS 38,000 52,000 +36.8%
P99 延迟 0.8ms 0.5ms -37.5%

鲲鹏的多核 + NUMA 架构对 Redis 这种单线程多实例场景特别友好,同样的 DCS 规格下 QPS 高出 30%+。

七、接入前后压测对比

场景 接入前 接入后 变化
下单接口 P99 420ms 86ms -79.5%
下单接口 TPS 280 1,850 +561%
库存扣减 TPS 220 3,500 +1491%
大促下单成功率 62% 99.8% +37.8pp
短信故障影响面 全链路挂 仅通知延迟 隔离
DB 读 QPS 8,500 1,200 -85.9%
超卖事件 大促偶发 0 清零

八、总结与展望

8.1 实践总结

引入 DMS Kafka + DCS Redis 后,订单中心从"同步串行"变成"异步解耦 + 缓存前置",三个最大收益:

  1. 下单主链路提速 5 倍:用户感知的耗时只剩预扣 + 创单 + 发消息
  2. 故障隔离:短信/积分等非核心链路挂了不影响下单
  3. 并发能力提升 15 倍:Redis Lua 预扣 + 异步 DB 扣减,扛住大促流量

三个最值得的决策:

  • 用 Outbox 模式保证消息可靠投递,比"先发消息再写库"稳得多
  • 用 Redisson 替代手写 SETNX,看门狗自动续期避免锁提前释放
  • 用 DRS 订阅 Binlog 刷缓存,比"先更 DB 再删缓存"一致性更强

8.2 后续规划

  • 接入 DMS for RocketMQ 的延迟消息实现"订单 30 分钟未支付自动取消",替代当前的定时任务扫表
  • 试点 DCS for Redis 的 RedisSearch 模块,实现订单多维度实时查询,替代 ES
  • 探索 华为云 GeminiDB(兼容 Redis 协议的分布式 KV)应对单集群容量瓶颈

互动引导:你的项目里缓存一致性用的是哪种方案?先删缓存再更 DB、先更 DB 再删缓存、还是 Binlog 订阅?遇到过哪些不一致的坑?欢迎评论区交流。

下一篇拆解 华为云 CSE + Spring Cloud Huawei 微服务治理,解决多服务间的注册发现、熔断降级、分布式事务问题,敬请关注。

【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。