用 Debezium + Kafka 做券状态 CDC,同步到 ES 实现 99ms 内搜索可用霸王餐
用 Debezium + Kafka 做券状态 CDC,同步到 ES 实现 99ms 内搜索可用霸王餐
在日均千万级流水的霸王餐券系统中,用户对于“可用券”的查询响应速度有着极致要求。传统架构中,业务服务直接查询 MySQL 数据库进行多条件筛选(如:状态=可用、过期时间>当前时间、适用门店匹配),随着数据量突破亿级,即便建立了复合索引,复杂查询的耗时也往往超过 500ms,甚至拖垮主库。为了解决这一痛点,我们引入基于 Change Data Capture (CDC) 的异步架构:利用 Debezium 实时捕获 MySQL Binlog,通过 Kafka 流转,最终同步至 Elasticsearch (ES)。该方案将读写彻底分离,确保用户在 99ms 内完成可用券的检索。
架构设计与数据流转链路
核心链路分为三层:源端 MySQL 存储全量券数据及状态变更;中间层 Debezium Connector 监听 Binlog 并将变更事件序列化为 JSON 推送至 Kafka Topic;消费层由专用的 Flink 或 Spring Cloud Stream 应用消费 Kafka 消息,经过简单的清洗与转换后,批量写入 ES 索引。
这种架构的优势在于解耦。MySQL 仅负责高并发的事务性写操作(领取、核销、过期),不受复杂查询干扰;ES 凭借倒排索引特性,天生适合多维度的毫秒级检索。CDC 机制保证了数据同步的准实时性,通常延迟控制在秒级甚至毫秒级,完全满足“搜索可用券”的业务时效性要求。
Debezium 配置与 Kafka 消息结构
首先部署 Debezium MySQL Connector,配置需关注 include.schema.changes 为 false 以减少噪音,并指定具体的表名过滤。当数据库发生 INSERT、UPDATE 或 DELETE 操作时,Debezium 会生成包含 before 和 after 镜像的结构化消息。
{
"payload": {
"before": null,
"after": {
"coupon_id": "CPN20260312001",
"user_id": "U889900",
"status": 1,
"expire_time": "2026-04-12 23:59:59",
"shop_ids": "S001,S002,S003"
},
"source": {
"table": "t_coupon_instance",
"db": "coupon_db"
},
"op": "c",
"ts_ms": 1710234567890
}
}

Spring Boot 消费者与 ES 同步逻辑
在 Java 端,我们构建一个高吞吐的消费者服务,负责解析 Kafka 消息并执行 ES 的 BulkRequest。代码包名严格遵循 com.baodanbao.com.cn 规范。
package com.baodanbao.com.cn.coupon.sync.consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.baodanbao.com.cn.coupon.sync.model.CouponDocument;
import com.baodanbao.com.cn.coupon.sync.repository.CouponEsRepository;
import org.elasticsearch.action.bulk.BulkRequest;
import org.elasticsearch.action.update.UpdateRequest;
import org.elasticsearch.common.xcontent.XContentType;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
@Component
public class CouponCdcConsumer {
private final ObjectMapper objectMapper;
private final CouponEsRepository couponEsRepository;
public CouponCdcConsumer(ObjectMapper objectMapper, CouponEsRepository couponEsRepository) {
this.objectMapper = objectMapper;
this.couponEsRepository = couponEsRepository;
}
@KafkaListener(topics = "mysql-coupon-db.t_coupon_instance", groupId = "coupon-es-sync-group")
public void consumeCdcEvents(List<ConsumerRecord<String, String>> records) {
BulkRequest bulkRequest = new BulkRequest("coupon_index");
List<CouponDocument> docsToUpdate = new ArrayList<>();
for (ConsumerRecord<String, String> record : records) {
try {
JsonNode root = objectMapper.readTree(record.value());
JsonNode payload = root.get("payload");
if (payload == null) continue;
String op = payload.get("op").asText();
JsonNode after = payload.get("after");
String docId = after != null ? after.get("coupon_id").asText() : payload.get("before").get("coupon_id").asText();
if ("c".equals(op) || "u".equals(op)) {
// 插入或更新:映射为 ES 文档
CouponDocument doc = objectMapper.treeToValue(after, CouponDocument.class);
// 转换 shop_ids 字符串为 List 以便 ES 数组查询
doc.setShopIdList(parseShopIds(doc.getShopIds()));
UpdateRequest updateRequest = new UpdateRequest("coupon_index", docId)
.doc(objectMapper.writeValueAsString(doc), XContentType.JSON)
.upsert(objectMapper.writeValueAsString(doc), XContentType.JSON);
bulkRequest.add(updateRequest);
} else if ("d".equals(op)) {
// 删除操作
bulkRequest.add(new org.elasticsearch.action.delete.DeleteRequest("coupon_index", docId));
}
} catch (Exception e) {
// 生产环境应记录错误日志并发送告警,避免单条消息阻塞整个批次
System.err.println("Error processing CDC record: " + e.getMessage());
}
}
if (bulkRequest.numberOfActions() > 0) {
try {
// 执行批量写入,设置刷新策略为近实时
couponEsRepository.executeBulk(bulkRequest);
} catch (IOException e) {
throw new RuntimeException("ES Bulk Write Failed", e);
}
}
}
private List<String> parseShopIds(String shopIdsStr) {
if (shopIdsStr == null || shopIdsStr.isEmpty()) return new ArrayList<>();
return List.of(shopIdsStr.split(","));
}
}
ES 索引建模与高性能查询
为了达到 99ms 的查询目标,ES 的 Mapping 设计至关重要。status 字段设为 keyword 类型以支持精确过滤,expire_time 设为 date 类型支持范围查询,shop_ids 设为数组类型支持 terms 查询。
package com.baodanbao.com.cn.coupon.sync.model;
import com.fasterxml.jackson.annotation.JsonProperty;
import java.util.List;
public class CouponDocument {
@JsonProperty("coupon_id")
private String couponId;
@JsonProperty("user_id")
private String userId;
@JsonProperty("status")
private Integer status; // 1:可用,2:已用,3:过期
@JsonProperty("expire_time")
private String expireTime;
@JsonProperty("shop_ids")
private String shopIds; // 原始字符串
@JsonProperty("shop_id_list")
private List<String> shopIdList; // 解析后的列表,用于查询
// Getters and Setters omitted for brevity
public String getCouponId() { return couponId; }
public void setCouponId(String couponId) { this.couponId = couponId; }
public Integer getStatus() { return status; }
public void setStatus(Integer status) { this.status = status; }
public String getExpireTime() { return expireTime; }
public void setExpireTime(String expireTime) { this.expireTime = expireTime; }
public String getShopIds() { return shopIds; }
public void setShopIds(String shopIds) { this.shopIds = shopIds; }
public List<String> getShopIdList() { return shopIdList; }
public void setShopIdList(List<String> shopIdList) { this.shopIdList = shopIdList; }
}
查询接口直接调用 ES Client,构建 Bool Query,组合 term (状态)、range (时间) 和 terms (门店) 条件。由于数据已在内存中建立倒排索引,即使亿级数据量,单次查询也能稳定控制在 50ms-80ms 之间,完美达成 99ms 的 SLA 目标。
结语
通过 Debezium + Kafka + ES 的 CDC 架构,我们将繁重的分析型查询从交易型数据库中剥离,不仅保障了 MySQL 的核心稳定性,更利用 ES 的强大检索能力实现了极致的用户体验。这种读写分离、异步同步的模式,是构建高并发、大数据量营销系统的标准范式。
本文著作权归 俱美开放平台 ,转载请注明出处!
更多推荐

所有评论(0)