多来买:商品上下架同步与一致性实战
多来买:商品上下架同步与一致性实战
事实来源:当前 product-service、search-service、common-mq 的上下架、Feign、ES 与 MQ 源码
证据状态:代码已确认,未启动 Nacos、Redis、RocketMQ 或 Elasticsearch。
阅读目标:理解商品状态如何跨 MySQL、Bloom、Feign、ES 和 MQ 传播,以及为什么当前双路径不能称为可靠最终一致性。
1. 业务问题
商品上架后需要同时满足:
- 商品库中的
is_sale变为上架; - 详情 Bloom 能识别这个 SKU;
- 搜索服务构建 ES 文档;
- 后续搜索可以检索到商品。
商品下架后需要撤销搜索可见性和热度状态。
当前项目同时保留同步 Feign 和 RocketMQ 两条通知路径,但它们的执行顺序、消费者状态和异常处理并没有形成完整补偿闭环。
2. 涉及组件
| 组件 | 当前职责 |
|---|---|
AdminSkuController | 提供上架、下架后台入口 |
SkuServiceImpl | 更新 MySQL、Bloom,并调用 Feign 和 MQ |
SearchApiClient | 商品服务调用搜索服务 |
SearchController | 接收上下架内部请求 |
SearchServiceImpl | 构建、保存或删除 Goods 文档 |
BaseProducer | 发送上下架 Topic |
MqOnsaleConsumer | 消费上架消息,当前缺少启动调用 |
MqOffSaleConsumer | 消费下架消息并启动消费者 |
3. 完整时序
4. 后台入口
源码:duolaimall-product/product-service/src/main/java/com/cskaoyan/mall/product/controller/AdminSkuController.java,onSale 与 cancelSale。
@GetMapping("admin/product/onSale/{skuId}")public Result onSale(@PathVariable Long skuId) { skuService.onSale(skuId); return Result.ok();}
@GetMapping("admin/product/cancelSale/{skuId}")public Result cancelSale(@PathVariable Long skuId) { skuService.offSale(skuId); return Result.ok();}两个 GET 接口都会改变服务端状态。
GET 通常应具有安全、可重复读取语义;用它执行写操作会引入缓存、预取和误触发风险,这是当前接口设计边界。
5. 商品上架顺序
源码:duolaimall-product/product-service/src/main/java/com/cskaoyan/mall/product/service/impl/SkuServiceImpl.java,SkuServiceImpl#onSale。
public void onSale(Long skuId) { SkuInfo skuInfo = new SkuInfo(); skuInfo.setIsSale(1); skuInfo.setId(skuId); skuInfoMapper.updateById(skuInfo);
RBloomFilter<Long> bloomFilter = redissonClient.getBloomFilter(RedisConst.SKU_BLOOM_FILTER); bloomFilter.add(skuId);
searchApiClient.upperGoods(skuId);
baseProducer.sendMessage( MqTopicConst.PRODUCT_ONSALE_TOPIC, skuId);}严格顺序:
MySQL is_sale=1 -> Bloom add -> 同步 Feign 上架 -> 发送上架 MQ当前没有检查:
- 数据库更新行数;
- Bloom
add的结果; - Feign 返回的业务状态;
- MQ 生产者返回的 Boolean。
6. 商品下架顺序
源码同上:SkuServiceImpl#offSale。
LambdaUpdateWrapper<SkuInfo> wrapper = new LambdaUpdateWrapper<>();wrapper.eq(SkuInfo::getId, skuId) .set(SkuInfo::getIsSale, 0);skuInfoMapper.update(null, wrapper);
searchApiClient.lowerGoods(skuId);
baseProducer.sendMessage( MqTopicConst.PRODUCT_OFFSALE_TOPIC, skuId);搜索服务的下架实现:
goodsRepository.deleteById(skuId);
RScoredSortedSet<Object> scores = redissonClient.getScoredSortedSet(RedisConst.HOT_SCORE);scores.remove(skuId);下架没有从普通 Bloom Filter 删除 SKU。
这不表示下架失败:Bloom 只是详情前置过滤;真正是否允许访问仍应由业务事实决定。当前详情代码是否严格阻止已下架 SKU,需要结合商品查询口径判断。
7. 搜索服务怎样取得商品数据
搜索服务不直连商品数据库,而是调用商品内部 API。
源码:duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/client/ProductApiClient.java,ProductApiClient。
@FeignClient("service-product")public interface ProductApiClient { @GetMapping("/api/product/inner/getSkuInfo/{skuId}") SkuInfoDTO getSkuInfo(@PathVariable Long skuId);
@GetMapping("/api/product/inner/getCategoryView/{categoryId}") CategoryHierarchyDTO getCategoryView( @PathVariable("categoryId") Long categoryId);
@GetMapping("/api/product/inner/getAttrList/{skuId}") List<PlatformAttributeInfoDTO> getAttrList( @PathVariable Long skuId);
@GetMapping("/api/product/inner/getTrademark/{tmId}") TrademarkDTO getTrademark(@PathVariable Long tmId);}路径变量的真实注解名称以源码为准;表格表达的是四类数据职责。
构建文档所需数据:
| API | Goods 字段 |
|---|---|
| SKU | ID、标题、图、价格、品牌 ID、分类 ID |
| 品牌 | 品牌名、Logo |
| 分类 | 一到三级 ID 与名称 |
| 平台属性 | Nested SearchAttr 列表 |
8. 同步 Feign 上架入口
源码:duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/controller/SearchController.java,upperGoods。
@GetMapping("/api/list/inner/upperGoods/{skuId}")public Result upperGoods(@PathVariable Long skuId) { searchService.upperGoodsAsync(skuId); return Result.ok();}方法名中的 Async 不等于 HTTP 立即返回,因为内部最终会调用 allOf(...).join()。
9. 构建 Goods 根字段
源码:duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/service/impl/SearchServiceImpl.java,upperGoodsAsync。
Goods goods = new Goods();ExecutorService executorService = Executors.newFixedThreadPool(8);
CompletableFuture<SkuInfoDTO> cf1 = CompletableFuture.supplyAsync(() -> { SkuInfoDTO skuInfo = productApiClient.getSkuInfo(skuId); if (skuInfo == null) { throw new RuntimeException( "获取 SKU 基本信息失败,skuId=" + skuId); }
goods.setId(skuInfo.getId()); goods.setTitle(skuInfo.getSkuName()); goods.setDefaultImg(skuInfo.getSkuDefaultImg()); goods.setPrice(skuInfo.getPrice().doubleValue()); return skuInfo; }, executorService);ES 文档 ID 直接使用 SKU ID,同一 SKU 再次 save 会定位同一文档标识。
价格转换为 Double 只服务搜索展示和排序,不能替代订单结算中的精确 BigDecimal 价格。
10. 品牌与分类依赖 SKU
源码同上:SearchServiceImpl#upperGoodsAsync。
CompletableFuture<Void> brandFuture = cf1.thenAcceptAsync(skuInfo -> { try { TrademarkDTO trademark = productApiClient.getTrademark(skuInfo.getTmId()); if (trademark != null) { goods.setTmId(skuInfo.getTmId()); goods.setTmName(trademark.getTmName()); goods.setTmLogoUrl(trademark.getLogoUrl()); } } catch (Exception e) { log.error("获取品牌信息失败 skuId={}", skuId, e); } }, executorService);品牌和分类分支捕获各自 Feign 异常,不继续传播,因此 ES 可能保存缺少部分字段的文档。
11. 平台属性转换
源码同上:SearchServiceImpl#upperGoodsAsync。
List<SearchAttr> searchAttrs = attrList.stream() .filter(Objects::nonNull) .map(attr -> { List<PlatformAttributeValueDTO> values = attr.getAttrValueList(); if (values == null || values.isEmpty()) { return null; }
PlatformAttributeValueDTO value = values.get(0); if (value == null || value.getValueName() == null) { return null; }
SearchAttr searchAttr = new SearchAttr(); searchAttr.setAttrId(attr.getId()); searchAttr.setAttrName(attr.getAttrName()); searchAttr.setAttrValue(value.getValueName()); return searchAttr; }) .filter(Objects::nonNull) .collect(Collectors.toList());每个平台属性只取第一个值,隐含“一项属性对一个确定 SKU 只有一个实际值”的假设。
12. ES 保存异常语义
源码同上:SearchServiceImpl#upperGoodsAsync。
CompletableFuture.allOf(cf2, cf3, cf4).join();
try { goodsRepository.save(goods); log.info("商品 [{}] 上架成功并保存至 Elasticsearch", skuId);} catch (Exception e) { log.error("商品 [{}] 保存到 ES 失败", skuId, e);}两类异常处理不同:
- SKU 根查询失败时,
join()可能向上抛出; - ES
save失败时,只记录日志,不继续抛出。
因此 Search Controller 返回 Result.ok() 不能证明文档已经存在。
upperGoodsAsync 的固定线程池也没有调用 shutdown。
13. 为什么 Feign 失败时 MQ 不一定兜底
Feign 在 MQ 发送语句之前,onSale 和 offSale 没有捕获 Feign 传输异常。
另一方面,如果 HTTP 调用成功,而搜索服务内部 ES 保存失败后仍返回成功,商品服务会继续发送 MQ。
MQ 能否成为第二次尝试,还要看消费者是否启动和消费逻辑是否可靠。
14. MQ 生产者语义
通用生产者把 skuId 序列化为 JSON,Topic 来自:
MqTopicConst.PRODUCT_ONSALE_TOPIC;MqTopicConst.PRODUCT_OFFSALE_TOPIC。
BaseProducer#sendMessage 成功返回 true,失败或异常返回 false。
SkuServiceImpl 没有使用返回值,所以 MQ 发送失败不会通过业务返回值改变后台接口结果。
15. 上架消费者没有启动
源码:duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/mq/MqOnsaleConsumer.java,init。敏感地址不在文档中回显。
defaultMQPushConsumer.subscribe( MqTopicConst.PRODUCT_ONSALE_TOPIC, "*");
defaultMQPushConsumer.setMessageListener( new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage( List<MessageExt> msgs, ConsumeConcurrentlyContext context) { Long skuId = JSON.parseObject( new String(msgs.get(0).getBody()), Long.class); searchService.upperGoods(skuId); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } });
// 文件末尾没有 defaultMQPushConsumer.start()代码创建消费者、订阅 Topic 并注册监听器,但缺少启动调用。
因此只能说“存在上架 MQ 消费代码”,不能说“上架消息当前能够被消费”。
16. 下架消费者已经调用 start
源码:duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/mq/MqOffSaleConsumer.java,init。
defaultMQPushConsumer.subscribe( MqTopicConst.PRODUCT_OFFSALE_TOPIC, "*");
defaultMQPushConsumer.setMessageListener( new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage( List<MessageExt> msgs, ConsumeConcurrentlyContext context) { Long skuId = JSON.parseObject( new String(msgs.get(0).getBody()), Long.class); searchService.lowerGoods(skuId); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } });
defaultMQPushConsumer.start();上、下架消费者还使用了不同的硬编码 MQ 地址,且没有统一复用生产者外部配置;本文只记录风险,不回显地址。
17. 两条上架实现的差异
| 路径 | 调用方法 | 异常特点 |
|---|---|---|
| 同步 Feign | upperGoodsAsync | 品牌、分类、属性局部捕获;ES save 吞异常 |
| MQ 消费 | upperGoods | 同步顺序读取,异常传播方式不同 |
两者都按同一 skuId 保存文档,但没有商品版本号或更新时间判断新旧。
若两条路径都执行,最后一次保存会覆盖同 ID 文档;不能保证最后执行的一定是更新数据。
18. 跨存储一致性边界
MySQL 更新 is_sale -> Redis Bloom -> Feign -> 搜索服务反向 Feign 取商品数据 -> Elasticsearch -> RocketMQ这些动作跨数据库、Redis、HTTP、ES 和 MQ,商品服务本地事务无法让它们原子提交。
onSale、offSale 本身也没有声明本地事务。
19. 失败矩阵
| 失败位置 | 可能留下的状态 | 自动补偿证据 |
|---|---|---|
| MySQL 更新未命中 | 后续仍可能继续 | 未检查更新行数 |
| Bloom 添加异常 | MySQL 已上架 | 未看到补偿 |
| Feign 传输异常 | MySQL 已变更,MQ 尚未发送 | 无 |
| ES save 异常 | 搜索服务可能仍返回成功 | 上架消费者未启动 |
| MQ 发送失败 | Feign 可能已完成 | 返回值被忽略 |
| 重复消费 | 同一文档重复 save/delete | 无显式幂等记录 |
| 旧消息后到 | 旧状态可能覆盖新状态 | 无版本比较 |
20. 为什么不能称为可靠最终一致性
可靠最终一致性通常还需要:
- 业务状态与待发送事件避免一边成功、一边丢失;
- 发送失败可重试;
- 消费失败返回正确重试状态;
- 重复消息幂等;
- 事件带版本,防止旧状态覆盖新状态;
- MySQL 与 ES 定期对账;
- 能按 SKU 重建索引。
当前代码没有形成这些完整闭环。
准确表述是:项目实现了 Feign 与 RocketMQ 两条上下架通知路径,用于练习搜索索引同步;消费者启动、失败处理、幂等、版本和对账仍不完整。
21. 已确认限制与候选风险
21.1 已确认限制
- 上架 MQ 消费者没有
start()。 - 上、下架消费者 MQ 地址没有统一配置。
- Feign 发生传输异常时,后面的 MQ 不会执行。
- ES 保存异常被捕获,HTTP 仍可能返回成功。
- MySQL、Bloom、ES 和 MQ 不在同一事务。
- MQ 返回值未被调用方使用。
- 两条文档写路径没有版本控制。
21.2 候选风险
| 风险 | 触发机制 | 缺失证据 |
|---|---|---|
| MySQL 已上架但搜索不可见 | Feign 或 ES 失败,上架消费未启动 | 未故障注入 |
| 部分字段文档 | 品牌、分类、属性失败后仍 save | 未读取真实 ES 文档 |
| 旧文档覆盖新文档 | 双路径无版本并发保存 | 未并发验证 |
| 缓存旧数据进入 ES | 搜索反向调用的商品 API 可能命中缓存 | 未连接 Redis 验证 |
| 下架后详情仍可进入后端查询 | Bloom 不删除 SKU | 未运行接口验证 |
22. 演进方向(非当前实现)
可以评估:
- 商品本地事务同时写业务状态和 outbox 事件;
- 独立任务可靠投递 outbox,成功后标记;
- 只保留一条权威索引写路径;
- 消费端按
skuId + 商品版本幂等更新; - 失败消息重试、死信和告警;
- 定时比较 MySQL 上架集合与 ES 文档集合;
- 提供按 SKU 和批次重建索引能力;
- 统一 MQ 地址和客户端配置。
这些建议不代表当前代码已经实现。
23. 关键源码导航
| 阅读问题 | 项目相对路径 | 类/方法 |
|---|---|---|
| 后台入口 | duolaimall-product/product-service/src/main/java/com/cskaoyan/mall/product/controller/AdminSkuController.java | onSale、cancelSale |
| 商品状态编排 | duolaimall-product/product-service/src/main/java/com/cskaoyan/mall/product/service/impl/SkuServiceImpl.java | onSale、offSale |
| 商品调搜索 | duolaimall-product/product-service/src/main/java/com/cskaoyan/mall/product/client/SearchApiClient.java | upperGoods、lowerGoods |
| 搜索入口 | duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/controller/SearchController.java | upperGoods、lowerGoods |
| 搜索文档写入 | duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/service/impl/SearchServiceImpl.java | upperGoodsAsync、upperGoods、lowerGoods |
| 商品内部 API 客户端 | duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/client/ProductApiClient.java | ProductApiClient |
| 通用 MQ 生产者 | duolaimall-common/common-mq/src/main/java/com/cskaoyan/mall/mq/producer/BaseProducer.java | sendMessage |
| 上架消费者 | duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/mq/MqOnsaleConsumer.java | init |
| 下架消费者 | duolaimall-search/search-service/src/main/java/com/cskaoyan/mall/search/mq/MqOffSaleConsumer.java | init |
24. 短复习点
- 上架顺序是 MySQL、Bloom、Feign、MQ,Feign 异常会阻止 MQ 发送。
- 搜索服务反向 Feign 获取商品数据,再构建 SKU 级 ES 文档。
- 上架消费者缺少
start(),不能把 MQ 当作当前可靠兜底。 - ES save 吞异常,使 HTTP 成功与索引成功不等价。
- 当前只能说实现了双通知路径,不能说已经实现可靠最终一致性。
25. 一句话总结
多来买商品上下架展示了 MySQL 状态、Redis Bloom、Feign、Elasticsearch 与 RocketMQ 的跨服务协作,但双路径缺少启动、幂等、版本、补偿和对账闭环,因此应定位为同步思路练习,而不是生产级最终一致性方案。
文章分享
如果这篇文章对你有帮助,欢迎分享给更多人!