削峰填谷与状态共享:华为云 DMS Kafka + DCS Redis 在 SpringBoot 中的实战
摘要: 本文复盘订单中心接入华为云 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。大促时并发一上来,三个问题同时爆雷:

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

核心收益:用户感知的"下单"只剩 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 订单中心中间件架构

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 再删缓存"在并发读写时仍会出现不一致:

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

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 后,订单中心从"同步串行"变成"异步解耦 + 缓存前置",三个最大收益:
- 下单主链路提速 5 倍:用户感知的耗时只剩预扣 + 创单 + 发消息
- 故障隔离:短信/积分等非核心链路挂了不影响下单
- 并发能力提升 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 微服务治理,解决多服务间的注册发现、熔断降级、分布式事务问题,敬请关注。
- 点赞
- 收藏
- 关注作者
评论(0)