Knowledge note
正在加载知识笔记
正在加载知识笔记
Knowledge note
模块定位:DDD 核心业务领域层
包路径:cn.bugstack.domain
依赖关系:仅依赖group-buy-market-types+xfg-wrench框架,被 app/infrastructure 依赖
源码文件:约 70 个 Java 类
group-buy-market-domain 是整个拼团营销系统的核心业务层。如果说 types 模块是"词汇表",那 domain 模块就是"造句"的地方——它不碰数据库、不碰 HTTP 接口、不碰 MQ,只定义业务做什么、怎么做、在什么条件下做。
┌──────────────────────────────────────────────┐
│ api / trigger ← 对外暴露的接口层 │
│ app ← 应用编排层 │
├──────────────────────────────────────────────┤
│ ★ domain ★ ← 你现在的位置 │
│ ├── activity ← 试算/报价子域 │
│ ├── trade ← 交易/锁单/结算/退单 │
│ └── tag ← 人群标签子域 │
├──────────────────────────────────────────────┤
│ infrastructure ← 数据库/Redis/MQ 实现 │
│ types ← 公共类型/枚举/异常 │
└──────────────────────────────────────────────┘
核心依赖规则:
这称为依赖倒置——domain 定义接口(IActivityRepository、ITradeRepository、ITradePort),infrastructure 提供实现。
cn.bugstack.domain/
├── activity/ # ★ 试算子域
│ ├── adapter/
│ │ └── repository/
│ │ └── IActivityRepository.java # 仓储接口(10个方法)
│ ├── model/
│ │ ├── entity/
│ │ │ ├── MarketProductEntity.java # 营销商品实体(请求入口)
│ │ │ ├── TrialBalanceEntity.java # 试算结果实体(返回出口)
│ │ │ └── UserGroupBuyOrderDetailEntity.java
│ │ └── valobj/
│ │ ├── DiscountTypeEnum.java # 折扣类型(BASE/TAG)
│ │ ├── GroupBuyActivityDiscountVO.java # 拼团活动配置VO
│ │ ├── SCSkuActivityVO.java # 渠道商品关联
│ │ ├── SkuVO.java # 商品信息
│ │ ├── TagScopeEnumVO.java # 标签范围枚举
│ │ └── TeamStatisticVO.java # 队伍统计
│ └── service/
│ ├── IIndexGroupBuyMarketService.java # 首页营销服务接口
│ ├── IndexGroupBuyMarketServiceImpl.java # 实现
│ ├── discount/ # 折扣计算子域
│ │ ├── IDiscountCalculateService.java
│ │ ├── AbstractDiscountCalculateService.java # 模板方法
│ │ └── impl/
│ │ ├── ZJCalculateService.java # 直减(@Service("ZJ"))
│ │ ├── MJCalculateService.java # 满减(@Service("MJ"))
│ │ ├── ZKCalculateService.java # 折扣(@Service("ZK"))
│ │ └── NCalculateService.java # N元购(@Service("N"))
│ └── trial/
│ ├── AbstractGroupBuyMarketSupport.java # 决策树抽象基类
│ ├── factory/
│ │ └── DefaultActivityStrategyFactory.java # 策略工厂 + DynamicContext
│ ├── node/
│ │ ├── RootNode.java # 根节点(参数校验 → 路由)
│ │ ├── SwitchNode.java # 开关节点(降级/切量判断)
│ │ ├── MarketNode.java # 营销节点(多线程加载 + 折扣计算)
│ │ ├── MarketNode2CompletableFuture.java # 多线程案例(CompletableFuture版)
│ │ ├── TagNode.java # 人群标签节点
│ │ ├── EndNode.java # 正常结束节点
│ │ └── ErrorNode.java # 异常结束节点
│ └── thread/
│ ├── QueryGroupBuyActivityDiscountVOThreadTask.java # 异步查询活动配置
│ └── QuerySkuVOFromDBThreadTask.java # 异步查询商品
│
├── trade/ # ★ 交易子域
│ ├── adapter/
│ │ ├── port/
│ │ │ └── ITradePort.java # 外部端口接口(HTTP回调)
│ │ └── repository/
│ │ └── ITradeRepository.java # 仓储接口(20+方法)
│ ├── model/
│ │ ├── aggregate/ # 聚合根 ★
│ │ │ ├── GroupBuyOrderAggregate.java # 锁单聚合
│ │ │ ├── GroupBuyTeamSettlementAggregate.java # 结算聚合
│ │ │ └── GroupBuyRefundAggregate.java # 退单聚合
│ │ ├── entity/ # 实体
│ │ │ ├── UserEntity.java
│ │ │ ├── GroupBuyActivityEntity.java
│ │ │ ├── GroupBuyTeamEntity.java
│ │ │ ├── MarketPayOrderEntity.java
│ │ │ ├── NotifyTaskEntity.java
│ │ │ ├── PayActivityEntity.java
│ │ │ ├── PayDiscountEntity.java
│ │ │ ├── TradeLockRuleCommandEntity.java
│ │ │ ├── TradeLockRuleFilterBackEntity.java
│ │ │ ├── TradePaySuccessEntity.java
│ │ │ ├── TradePaySettlementEntity.java
│ │ │ ├── TradeRefundCommandEntity.java
│ │ │ ├── TradeRefundOrderEntity.java
│ │ │ ├── TradeRefundBehaviorEntity.java
│ │ │ ├── TradeSettlementRuleCommandEntity.java
│ │ │ └── TradeSettlementRuleFilterBackEntity.java
│ │ └── valobj/ # 值对象
│ │ ├── GroupBuyProgressVO.java
│ │ ├── NotifyConfigVO.java
│ │ ├── NotifyTypeEnumVO.java
│ │ ├── RefundTypeEnumVO.java # ★ 退单策略分发枚举
│ │ ├── TaskNotifyCategoryEnumVO.java
│ │ ├── TeamRefundSuccess.java
│ │ └── TradeOrderStatusEnumVO.java
│ └── service/
│ ├── ITradeLockOrderService.java
│ ├── ITradeSettlementOrderService.java
│ ├── ITradeRefundOrderService.java
│ ├── ITradeTaskService.java
│ ├── lock/
│ │ ├── TradeLockOrderService.java # 锁单服务
│ │ ├── factory/
│ │ │ └── TradeLockRuleFilterFactory.java # 锁单责任链工厂
│ │ └── filter/
│ │ ├── ActivityUsabilityRuleFilter.java # 活动可用性校验
│ │ ├── UserTakeLimitRuleFilter.java # 用户参与次数限制
│ │ └── TeamStockOccupyRuleFilter.java # 组队库存占用
│ ├── settlement/
│ │ ├── TradeSettlementOrderService.java # 结算服务
│ │ ├── factory/
│ │ │ └── TradeSettlementRuleFilterFactory.java
│ │ └── filter/
│ │ ├── SCRuleFilter.java # SC渠道黑名单
│ │ ├── OutTradeNoRuleFilter.java # 外部单号校验
│ │ ├── SettableRuleFilter.java # 可结算时间校验
│ │ └── EndRuleFilter.java # 结束节点
│ ├── refund/
│ │ ├── TradeRefundOrderService.java # 退单服务
│ │ ├── factory/
│ │ │ └── TradeRefundRuleFilterFactory.java
│ │ ├── filter/
│ │ │ ├── DataNodeFilter.java # 数据加载
│ │ │ ├── UniqueRefundNodeFilter.java # 重复退单检查
│ │ │ └── RefundOrderNodeFilter.java # 策略分发
│ │ └── business/
│ │ ├── IRefundOrderStrategy.java
│ │ ├── AbstractRefundOrderStrategy.java # 退单抽象基类
│ │ └── impl/
│ │ ├── Unpaid2RefundStrategy.java # 未支付退单
│ │ ├── Paid2RefundStrategy.java # 已支付未成团退单
│ │ └── PaidTeam2RefundStrategy.java # 已支付已成团退单
│ └── task/
│ └── TradeTaskService.java # 回调通知任务服务
│
└── tag/ # ★ 标签子域
├── adapter/
│ └── repository/
│ └── ITagRepository.java
├── model/
│ └── entity/
│ └── CrowdTagsJobEntity.java
└── service/
├── ITagService.java
└── TagService.java
┌──────────────────────┐
│ trigger 层 │
│ (HTTP / MQ Listener) │
└────┬────────┬────────┘
│ │
┌─────────────┘ └─────────────┐
▼ ▼
┌─────────────────────┐ ┌──────────────────────┐
│ activity 子域 │ │ trade 子域 │
│ │ │ │
│ 首页试算 (Trial) │ │ 锁单 (Lock) │
│ Root → Switch → │ │ 结算 (Settlement) │
│ Market → Tag → End │ │ 退单 (Refund) │
│ │ │ 回调通知 (Task) │
│ 折扣计算 (Discount) │ │ │
└────────┬────────────┘ └───────────┬───────────┘
│ │
└──────────────┬──────────────────────┘
│
┌────────────┴────────────┐
│ tag 子域 │
│ 人群标签处理 │
└─────────────────────────┘
│
▼
┌────────────────────────┐
│ infrastructure 层 │
│ (MySQL / Redis / MQ) │
└────────────────────────┘
三个子域的分工:
| 子域 | 核心职责 | 输入 | 输出 |
|---|---|---|---|
| activity | 判断商品能否参加拼团、计算优惠价格 | MarketProductEntity (userId, goodsId, source, channel) | TrialBalanceEntity (折扣价、可见/可参与) |
| trade | 锁单、支付结算、退单、回调通知 | 订单/支付/退款事件 | 订单状态变更、通知任务 |
| tag | 人群标签规则匹配 | tagId, batchId | 标签人群数据写入 |
一句话:用户打开商品详情页时,判断该商品有没有拼团活动、用户能不能参与、能省多少钱。
核心流程:
MarketProductEntity → [决策树] → TrialBalanceEntity
(userId, goodsId) (6个节点) (折扣价, 可见性, 可参与性)
一句话:用户下单后的完整交易生命周期——锁单确认 → 支付结算 → 成团/退单 → 异步通知。
核心流程:
下单锁单 → 支付结算 → 拼团成团/未成团 → 退单处理 → 回调通知
↓ ↓ ↓ ↓ ↓
Lock Settlement 团状态变更 Refund Notify
一句话:维护哪些用户属于某个人群标签,供试算和折扣计算时做过滤。
IIndexGroupBuyMarketService (activity/service/IIndexGroupBuyMarketService.java:15)
public interface IIndexGroupBuyMarketService {
// 核心方法:商品试算
TrialBalanceEntity indexMarketTrial(MarketProductEntity marketProductEntity);
// 查询进行中的拼团(我的 + 随机的)
List<UserGroupBuyOrderDetailEntity> queryInProgressUserGroupBuyOrderDetailList(
Long activityId, String userId, Integer ownerCount, Integer randomCount);
// 查询活动队伍统计
TeamStatisticVO queryTeamStatisticByActivityId(Long activityId);
}
第一个方法 indexMarketTrial 是 activity 子域最核心的入口,整个决策树就是为了处理这一次调用。
@Data @Builder
public class MarketProductEntity {
private Long activityId; // 活动ID(可为空,为空时通过source+channel+goodsId查询)
private String userId; // 用户ID
private String goodsId; // 商品ID
private String source; // 来源
private String channel; // 渠道
}
activityId 可为空——如果为空,MarketNode 内部会通过 source + channel + goodsId 去查 sc_sku_activity 表找到关联的活动ID。
@Data @Builder
public class TrialBalanceEntity {
private String goodsId; // 商品ID
private String goodsName; // 商品名称
private BigDecimal originalPrice; // 原始价格
private BigDecimal deductionPrice; // 折扣金额
private BigDecimal payPrice; // 支付金额(原价 - 折扣)
private Integer targetCount; // 拼团需要多少人
private Date startTime; // 活动开始时间
private Date endTime; // 活动结束时间
private Boolean isVisible; // 是否对用户可见拼团价
private Boolean isEnable; // 用户是否可参与拼团
private GroupBuyActivityDiscountVO groupBuyActivityDiscountVO; // 完整活动配置
}
关键字段解读:
payPrice = originalPrice - deductionPrice:用户实际要付的钱isVisible:该用户能不能看到拼团优惠价(人群过滤)isEnable:该用户能不能发起/参与拼团(人群过滤)这是 activity 子域最重要的值对象,承载活动配置的完整信息。
@Getter @Builder
public class GroupBuyActivityDiscountVO {
private Long activityId;
private String activityName;
private Integer groupType; // 0=自动成团, 1=达成目标拼团
private Integer takeLimitCount; // 每人最多参与次数
private Integer target; // 成团目标人数
private Integer validTime; // 有效时长(分钟)
private Integer status; // 活动状态
private Date startTime;
private Date endTime;
private String tagId; // 人群标签ID
private String tagScope; // 标签范围规则
private GroupBuyDiscount groupBuyDiscount; // 折扣详情
// 内嵌类:折扣配置
public static class GroupBuyDiscount {
private String discountName;
private String discountDesc;
private DiscountTypeEnum discountType; // BASE(0) / TAG(1)
private String marketPlan; // ZJ/MJ/ZK/N
private String marketExpr; // 表达式(如"100,10"表示满100减10)
private String tagId; // 限定人群标签
}
}
这是整个 domain 模块最核心的设计模式应用——使用 xfg-wrench 的 AbstractMultiThreadStrategyRouter 构建了一个树形责任链。
┌──────────┐
┌───→│ RootNode │ (参数校验)
│ └────┬─────┘
│ │ router()
│ ▼
│ ┌────────────┐
│ │ SwitchNode │ (降级/切量判断)
│ └─────┬──────┘
│ │ router()
│ ▼
│ ┌────────────┐
│ │ MarketNode │ (多线程加载数据 + 折扣计算)
│ └─────┬──────┘
│ │ router()─→ 无配置 → ErrorNode
│ ▼
│ ┌──────────┐
│ │ TagNode │ (人群标签过滤)
│ └────┬─────┘
│ │ router()
│ ▼
│ ┌──────────┐
│ │ EndNode │ (组装返回结果)
│ └──────────┘
│
└── ErrorNode (异常/无营销 → 抛异常)
树形结构而非链式结构——每个节点通过 get() 方法决定下一个节点是谁,doApply() 执行业务逻辑,router() 串联两者。
public abstract class AbstractGroupBuyMarketSupport
extends AbstractMultiThreadStrategyRouter<MarketProductEntity,
DefaultActivityStrategyFactory.DynamicContext, TrialBalanceEntity> {
protected long timeout = 5000;
@Resource
protected IActivityRepository repository;
@Override
protected void multiThread(...) {
// 默认空实现,子类可覆写实现多线程数据加载
}
}
关键继承链:
AbstractMultiThreadStrategyRouter(xfg-wrench)— 提供决策树路由框架 + 多线程支持AbstractGroupBuyMarketSupport(domain)— 注入 repository、设置超时每个节点的执行流程是固定的:
multiThread(request, ctx) → 异步预加载数据(可选)
↓
doApply(request, ctx) → 执行业务逻辑
↓
router(request, ctx) → 内部调 get() + next.apply()
↓
get(request, ctx) → 返回下一个节点
@Service
public class RootNode extends AbstractGroupBuyMarketSupport<...> {
@Resource
private SwitchNode switchNode;
@Override
protected TrialBalanceEntity doApply(MarketProductEntity requestParameter,
DefaultActivityStrategyFactory.DynamicContext dynamicContext) {
// 参数校验:userId、goodsId、source、channel 不能为空
if (StringUtils.isBlank(requestParameter.getUserId()) ...) {
throw new AppException(ResponseCode.ILLEGAL_PARAMETER);
}
return router(requestParameter, dynamicContext);
}
@Override
public StrategyHandler<...> get(...) {
return switchNode; // 固定路由到 SwitchNode
}
}
职责:参数前置校验。入口的第一道关卡。
@Service
public class SwitchNode extends AbstractGroupBuyMarketSupport<...> {
@Resource
private MarketNode marketNode;
@Override
public TrialBalanceEntity doApply(...) {
// 1. 降级开关
if (repository.downgradeSwitch()) {
throw new AppException(ResponseCode.E0003); // "拼团活动降级拦截"
}
// 2. 切量判断
if (!repository.cutRange(userId)) {
throw new AppException(ResponseCode.E0004); // "拼团活动切量拦截"
}
return router(requestParameter, dynamicContext);
}
@Override
public StrategyHandler<...> get(...) {
return marketNode; // 固定路由到 MarketNode
}
}
职责:灰度控制。在基础设施层通过 Redis/配置中心控制哪些用户能看到拼团活动。这是一个典型的**特性开关(Feature Flag)**模式。
这是整个决策树最复杂的节点。
@Service
public class MarketNode extends AbstractGroupBuyMarketSupport<...> {
@Resource
private ThreadPoolExecutor threadPoolExecutor;
@Resource
private Map<String, IDiscountCalculateService> discountCalculateServiceMap; // Spring 注入所有折扣策略
@Resource
private ErrorNode errorNode;
@Resource
private TagNode tagNode;
// 覆写 multiThread:多线程并行加载数据
@Override
protected void multiThread(MarketProductEntity requestParameter,
DefaultActivityStrategyFactory.DynamicContext dynamicContext) {
// 1. 异步查询活动配置
QueryGroupBuyActivityDiscountVOThreadTask task1 = ...;
FutureTask<GroupBuyActivityDiscountVO> future1 = new FutureTask<>(task1);
threadPoolExecutor.execute(future1);
// 2. 异步查询商品信息
QuerySkuVOFromDBThreadTask task2 = ...;
FutureTask<SkuVO> future2 = new FutureTask<>(task2);
threadPoolExecutor.execute(future2);
// 3. 等待结果写入 DynamicContext
dynamicContext.setGroupBuyActivityDiscountVO(future1.get(timeout, TimeUnit.MILLISECONDS));
dynamicContext.setSkuVO(future2.get(timeout, TimeUnit.MILLISECONDS));
}
@Override
public TrialBalanceEntity doApply(...) {
// 从 DynamicContext 获取数据
GroupBuyActivityDiscountVO discountVO = dynamicContext.getGroupBuyActivityDiscountVO();
SkuVO skuVO = dynamicContext.getSkuVO();
// 通过 marketPlan(ZJ/MJ/ZK/N)找到对应的折扣计算策略
IDiscountCalculateService service = discountCalculateServiceMap.get(groupBuyDiscount.getMarketPlan());
// 执行计算
BigDecimal payPrice = service.calculate(userId, skuVO.getOriginalPrice(), groupBuyDiscount);
dynamicContext.setDeductionPrice(originalPrice.subtract(payPrice));
dynamicContext.setPayPrice(payPrice);
return router(requestParameter, dynamicContext);
}
@Override
public StrategyHandler<...> get(...) {
if (null == dynamicContext.getGroupBuyActivityDiscountVO()
|| null == dynamicContext.getSkuVO()
|| null == dynamicContext.getDeductionPrice()) {
return errorNode; // 无配置 → 异常节点
}
return tagNode; // 正常 → 人群标签节点
}
}
设计亮点:
多线程并行加载:活动配置和商品信息是独立的,并行查询减少 RT。multiThread() 在 doApply() 之前由框架自动调用。
策略模式注入:Map<String, IDiscountCalculateService> 自动收集 Spring 中所有 @Service("ZJ")、@Service("MJ") 等 bean,通过 marketPlan 字段做策略分发。
MarketNode2CompletableFuture:作者还提供了 CompletableFuture 版本的实现(默认注释掉 @Service),展示了 FutureTask vs CompletableFuture 的用法对比。
@Service
public class TagNode extends AbstractGroupBuyMarketSupport<...> {
@Resource
private EndNode endNode;
@Override
protected TrialBalanceEntity doApply(...) {
GroupBuyActivityDiscountVO discountVO = dynamicContext.getGroupBuyActivityDiscountVO();
// tagId 为空 → 无人群限制 → 直接可见+可参与
if (StringUtils.isBlank(discountVO.getTagId())) {
dynamicContext.setVisible(true);
dynamicContext.setEnable(true);
return router(requestParameter, dynamicContext);
}
// 判断用户是否在人群范围内
boolean isWithin = repository.isTagCrowdRange(tagId, requestParameter.getUserId());
dynamicContext.setVisible(visible || isWithin);
dynamicContext.setEnable(enable || isWithin);
return router(requestParameter, dynamicContext);
}
}
逻辑说明:
visible / enable 来自 GroupBuyActivityDiscountVO.isVisible() / isEnable() 方法,它们根据 tagScope 字段解析:
tagScope = "1" → 可见性限制tagScope = "2" → 参与限制tagScope = "1,2" → 两者都限制isWithin:用户是否在人群标签中(查询 crowd_tags_detail 表)@Service
public class EndNode extends AbstractGroupBuyMarketSupport<...> {
@Override
public TrialBalanceEntity doApply(...) {
// 从 DynamicContext 组装最终结果
return TrialBalanceEntity.builder()
.goodsId(skuVO.getGoodsId())
.goodsName(skuVO.getGoodsName())
.originalPrice(skuVO.getOriginalPrice())
.deductionPrice(dynamicContext.getDeductionPrice())
.payPrice(dynamicContext.getPayPrice())
.targetCount(discountVO.getTarget())
.startTime(discountVO.getStartTime())
.endTime(discountVO.getEndTime())
.isVisible(dynamicContext.isVisible())
.isEnable(dynamicContext.isEnable())
.groupBuyActivityDiscountVO(discountVO)
.build();
}
}
职责:收口——把 DynamicContext 里攒好的所有数据组装成 TrialBalanceEntity 返回。
public static class DynamicContext {
private GroupBuyActivityDiscountVO groupBuyActivityDiscountVO;
private SkuVO skuVO;
private BigDecimal deductionPrice;
private BigDecimal payPrice;
private boolean visible;
private boolean enable;
}
DynamicContext 是决策树节点的共享数据载体。每个节点在处理时从 ctx 读数据、写数据,传给下游。类似 ThreadLocal 的用法,但它是显式传递的,线程安全。
QueryGroupBuyActivityDiscountVOThreadTask:实现 Callable<GroupBuyActivityDiscountVO>,在线程池中异步查询活动配置。activityId 为空时自动查 sc_sku_activity 关联表。
QuerySkuVOFromDBThreadTask:实现 Callable<SkuVO>,异步查询商品信息。
AbstractDiscountCalculateService 定义了计算流程的骨架:
public abstract class AbstractDiscountCalculateService implements IDiscountCalculateService {
@Override
public BigDecimal calculate(String userId, BigDecimal originalPrice,
GroupBuyActivityDiscountVO.GroupBuyDiscount groupBuyDiscount) {
// 1. 人群标签过滤(折扣层面的标签)
if (DiscountTypeEnum.TAG.equals(groupBuyDiscount.getDiscountType())){
if (!filterTagId(userId, groupBuyDiscount.getTagId())) {
return originalPrice; // 不在人群范围内 → 无折扣
}
}
// 2. 交给子类实现具体的计算逻辑
return doCalculate(originalPrice, groupBuyDiscount);
}
protected abstract BigDecimal doCalculate(...);
}
注意:这里是折扣层面的人群标签,与 TagNode 的活动层面的人群标签不同。一个控制"能不能打折",一个控制"能不能看到/参与活动"。
| 策略名 | marketPlan | marketExpr 示例 | 计算逻辑 |
|---|---|---|---|
ZJCalculateService | "ZJ" | "20" | 原价 - 20,最低 0.01 |
MJCalculateService | "MJ" | "100,10" | 满 100 减 10,不满足返原价 |
ZKCalculateService | "ZK" | "0.85" | 原价 × 0.85,四舍五入 |
NCalculateService | "N" | "9.9" | 直接返回 9.9 元(N元购) |
Spring 注入技巧:通过 @Service("ZJ") 指定 bean 名称,Map<String, IDiscountCalculateService> 自动收集所有实现,key 为 bean 名,value 为实例。这样新增折扣类型只需加一个类,无需修改工厂代码。
交易子域是整个系统的核心,处理从"用户下单"到"成团结算"的完整流程。锁单是第一步。
TradeLockOrderService.lockMarketPayOrder() (trade/service/lock/TradeLockOrderService.java:42)
public MarketPayOrderEntity lockMarketPayOrder(
UserEntity userEntity,
PayActivityEntity payActivityEntity,
PayDiscountEntity payDiscountEntity) {
// 1. 责任链过滤
TradeLockRuleFilterBackEntity back = tradeRuleFilter.apply(
TradeLockRuleCommandEntity.builder()
.activityId(...).userId(...).teamId(...).build(),
new DynamicContext());
// 2. 构建聚合对象
GroupBuyOrderAggregate aggregate = GroupBuyOrderAggregate.builder()
.userEntity(userEntity)
.payActivityEntity(payActivityEntity)
.payDiscountEntity(payDiscountEntity)
.userTakeOrderCount(back.getUserTakeOrderCount())
.build();
// 3. 持久化(通过 repository 接口,由 infrastructure 实现)
return repository.lockMarketPayOrder(aggregate);
}
ActivityUsabilityRuleFilter → UserTakeLimitRuleFilter → TeamStockOccupyRuleFilter
(活动可用) (用户次数限制) (组队库存占用)
组装代码:
@Bean("tradeRuleFilter")
public BusinessLinkedList<...> tradeRuleFilter(...) {
LinkArmory<...> linkArmory = new LinkArmory<>("交易规则过滤链",
activityUsabilityRuleFilter,
userTakeLimitRuleFilter,
teamStockOccupyRuleFilter);
return linkArmory.getLogicLink();
}
校验活动是否有效:
GroupBuyActivityEntityEFFECTIVE[startTime, endTime] 范围内AppException校验用户参与次数:
takeLimitCount 上限 → 抛 AppException(E0103)DynamicContext校验组队库存——锁单流程中最巧妙的设计。
// teamId 为空 → 首次开团 → 不做库存限制
if (StringUtils.isBlank(teamId)) {
return ...; // 直接放行
}
// Redis 抢占库存
boolean status = repository.occupyTeamStock(
teamStockKey, // "group_buy_market_team_stock_key_{activityId}_{teamId}"
recoveryTeamStockKey, // 同上 + "_recovery"
target, // 成团目标人数
validTime // 有效期
);
设计意图:通过 Redis 缓存拼团组队剩余库存,减少对数据库 group_buy_order_list 表的 count 查询压力。库存恢复 key 用于支付失败/超时回退。
@Data @Builder
public class GroupBuyOrderAggregate {
private UserEntity userEntity; // 用户
private PayActivityEntity payActivityEntity; // 支付活动
private PayDiscountEntity payDiscountEntity; // 支付优惠
private Integer userTakeOrderCount; // 已参与次数(用于构建唯一索引)
}
聚合根的作用:把"锁单"这个业务操作需要的所有实体打包成一个整体,由 infrastructure 层一次事务写入多张表(group_buy_order + group_buy_order_list)。
TradeSettlementOrderService.settlementMarketPayOrder() (trade/service/settlement/TradeSettlementOrderService.java:43)
public TradePaySettlementEntity settlementMarketPayOrder(
TradePaySuccessEntity tradePaySuccessEntity) {
// 1. 责任链过滤
TradeSettlementRuleFilterBackEntity back = tradeSettlementRuleFilter.apply(
TradeSettlementRuleCommandEntity.builder()
.source(...).channel(...).userId(...).outTradeNo(...).outTradeTime(...).build(),
new DynamicContext());
// 2. 构建聚合
GroupBuyTeamSettlementAggregate aggregate = GroupBuyTeamSettlementAggregate.builder()
.userEntity(...).groupBuyTeamEntity(...).tradePaySuccessEntity(...).build();
// 3. 持久化结算
NotifyTaskEntity notifyTaskEntity = repository.settlementMarketPayOrder(aggregate);
// 4. 异步回调通知(如有需要)
if (null != notifyTaskEntity) {
threadPoolExecutor.execute(() -> {
tradeTaskService.execNotifyJob(notifyTaskEntity);
});
}
return TradePaySettlementEntity.builder()...build();
}
SCRuleFilter → OutTradeNoRuleFilter → SettableRuleFilter → EndRuleFilter
(渠道黑名单) (外部单号校验) (可结算时间) (数据组装)
四个过滤器的职责:
| 过滤器 | 校验内容 | 不通过 → |
|---|---|---|
| SCRuleFilter | source+channel 是否在黑名单 | 抛 E0105 |
| OutTradeNoRuleFilter | 外部单号是否存在、是否已退单 | 抛 E0104 |
| SettableRuleFilter | 支付时间是否在拼团有效期内 | 抛 E0106 |
| EndRuleFilter | 从 DynamicContext 组装返回数据 | — |
@Data @Builder
public class GroupBuyTeamSettlementAggregate {
private UserEntity userEntity;
private GroupBuyTeamEntity groupBuyTeamEntity;
private TradePaySuccessEntity tradePaySuccessEntity;
}
结算聚合包含了结算所需的全部信息,infrastructure 层通过它更新 group_buy_order_list(锁单→完成)、group_buy_order(状态变更)等表。
退单是整个系统最复杂的流程,因为退单场景取决于拼团状态 + 支付状态的组合。
TradeRefundOrderService.refundOrder() (trade/service/refund/TradeRefundOrderService.java:46)
public TradeRefundBehaviorEntity refundOrder(
TradeRefundCommandEntity tradeRefundCommandEntity) {
return tradeRefundRuleFilter.apply(
tradeRefundCommandEntity, new DynamicContext());
}
DataNodeFilter → UniqueRefundNodeFilter → RefundOrderNodeFilter
(数据加载) (幂等性检查) (策略分发+执行)
加载数据到 DynamicContext:
MarketPayOrderEntity(通过 outTradeNo → 获知 teamId + orderId)GroupBuyTeamEntity(通过 teamId → 获知拼团状态)幂等性保证:
if (TradeOrderStatusEnumVO.CLOSE.equals(tradeOrderStatusEnumVO)) {
return TradeRefundBehaviorEntity.builder()
.tradeRefundBehaviorEnum(TradeRefundBehaviorEnum.REPEAT)
.build(); // 已退单,直接返回幂等结果
}
策略分发——退单流程的设计核心:
// 1. 根据拼团状态 + 交易状态 获取退单类型
RefundTypeEnumVO refundType = RefundTypeEnumVO.getRefundStrategy(
groupBuyOrderEnumVO, // 拼团队伍状态
tradeOrderStatusEnumVO // 交易订单状态
);
// 2. 获取对应的策略实现
IRefundOrderStrategy strategy = refundOrderStrategyMap.get(refundType.getStrategy());
// 3. 执行退单
strategy.refundOrder(...);
这是 domain 模块最出彩的设计:
public enum RefundTypeEnumVO {
UNPAID_UNLOCK("unpaid_unlock", "unpaid2RefundStrategy", "未支付,未成团") {
@Override
public boolean matches(GroupBuyOrderEnumVO g, TradeOrderStatusEnumVO t) {
return PROGRESS.equals(g) && CREATE.equals(t);
}
},
PAID_UNFORMED("paid_unformed", "paid2RefundStrategy", "已支付,未成团") {
@Override
public boolean matches(GroupBuyOrderEnumVO g, TradeOrderStatusEnumVO t) {
return PROGRESS.equals(g) && COMPLETE.equals(t);
}
},
PAID_FORMED("paid_formed", "paidTeam2RefundStrategy", "已支付,已成团") {
@Override
public boolean matches(GroupBuyOrderEnumVO g, TradeOrderStatusEnumVO t) {
return (COMPLETE.equals(g) || COMPLETE_FAIL.equals(g)) && COMPLETE.equals(t);
}
};
public abstract boolean matches(GroupBuyOrderEnumVO, TradeOrderStatusEnumVO);
public static RefundTypeEnumVO getRefundStrategy(GroupBuyOrderEnumVO g, TradeOrderStatusEnumVO t) {
return Arrays.stream(values())
.filter(refundType -> refundType.matches(g, t))
.findFirst()
.orElseThrow(...);
}
}
设计精妙之处:
matches() 方法)strategy 字段存储 Spring bean 名称,直接映射到 @Service("unpaid2RefundStrategy")拼团状态 GroupBuyOrderEnumVO | 交易状态 TradeOrderStatusEnumVO | 退单策略 | 含义 |
|---|---|---|---|
| PROGRESS (拼单中) | CREATE (已创建,未支付) | Unpaid2RefundStrategy | 锁单但未付款 |
| PROGRESS (拼单中) | COMPLETE (已支付) | Paid2RefundStrategy | 已付款但团没成 |
| COMPLETE/COMPLETE_FAIL | COMPLETE | PaidTeam2RefundStrategy | 团已成功,退已支付订单 |
@Override
public void refundOrder(TradeRefundOrderEntity tradeRefundOrderEntity) {
// 1. 调 repository:锁单量 -1、更新订单状态
NotifyTaskEntity notifyTask = repository.unpaid2Refund(
GroupBuyRefundAggregate.buildUnpaid2RefundAggregate(entity, -1));
// 2. 发送 MQ 消息 → 异步恢复 Redis 库存
sendRefundNotifyMessage(notifyTask, "未支付,未成团");
}
@Override
public void refundOrder(TradeRefundOrderEntity tradeRefundOrderEntity) {
// 1. 锁单量 -1、完成量 -1
NotifyTaskEntity notifyTask = repository.paid2Refund(
GroupBuyRefundAggregate.buildPaid2RefundAggregate(entity, -1, -1));
// 2. MQ → 异步恢复库存
sendRefundNotifyMessage(notifyTask, "已支付,未成团");
}
@Override
public void refundOrder(TradeRefundOrderEntity tradeRefundOrderEntity) {
// 最后一笔退单 → 整个团标记为 FAIL
GroupBuyOrderEnumVO status =
1 == completeCount ? GroupBuyOrderEnumVO.FAIL : GroupBuyOrderEnumVO.COMPLETE_FAIL;
// 退单
NotifyTaskEntity notifyTask = repository.paidTeam2Refund(
GroupBuyRefundAggregate.buildPaidTeam2RefundAggregate(entity, -1, -1, status));
sendRefundNotifyMessage(notifyTask, "已支付,已成团");
}
关键逻辑:已成团退单时,如果退的是最后一笔(completeCount == 1),队伍状态从 COMPLETE/COMPLETE_FAIL 变为 FAIL。
public abstract class AbstractRefundOrderStrategy implements IRefundOrderStrategy {
@Resource
protected ITradeRepository repository; // 仓储
@Resource
protected ITradeTaskService tradeTaskService; // 任务通知
@Resource
protected ThreadPoolExecutor threadPoolExecutor;
// 异步发送 MQ 回调
protected void sendRefundNotifyMessage(NotifyTaskEntity notifyTask, String refundType) {
threadPoolExecutor.execute(() -> {
tradeTaskService.execNotifyJob(notifyTask);
});
}
// 通用库存恢复
protected void doReverseStock(TeamRefundSuccess teamRefundSuccess, String refundType) {
String recoveryKey = TradeLockRuleFilterFactory.generateRecoveryTeamStockKey(
teamRefundSuccess.getActivityId(), teamRefundSuccess.getTeamId());
repository.refund2AddRecovery(recoveryKey, teamRefundSuccess.getOrderId());
}
}
| 聚合 | 使用场景 | 核心实体 |
|---|---|---|
GroupBuyOrderAggregate | 锁单 | User + PayActivity + PayDiscount |
GroupBuyTeamSettlementAggregate | 结算 | User + GroupBuyTeam + TradePaySuccess |
GroupBuyRefundAggregate | 退单 | TradeRefundOrder + GroupBuyProgress + 状态枚举 |
负责在拼团成团/失败/退单完成后,通过 HTTP/MQ 通知外部系统。
@Service
public class TradeTaskService implements ITradeTaskService {
// 定时任务:查询所有未执行的通知任务
public Map<String, Integer> execNotifyJob() {
List<NotifyTaskEntity> list = repository.queryUnExecutedNotifyTaskList();
return execNotifyJob(list);
}
// 指定 teamId 重跑
public Map<String, Integer> execNotifyJob(String teamId) {
...
}
// 单条通知(结算/退单后的实时回调)
public Map<String, Integer> execNotifyJob(NotifyTaskEntity notifyTaskEntity) {
...
}
private Map<String, Integer> execNotifyJob(List<NotifyTaskEntity> list) {
for (NotifyTaskEntity task : list) {
String response = port.groupBuyNotify(task); // HTTP/MQ 回调
if ("success".equals(response)) {
repository.updateNotifyTaskStatusSuccess(task); // 标记成功
} else if ("error".equals(response)) {
if (task.getNotifyCount() > 4) {
repository.updateNotifyTaskStatusError(task); // 重试超5次→失败
} else {
repository.updateNotifyTaskStatusRetry(task); // 重试次数+1
}
}
}
}
}
重试策略:最多重试 5 次(notifyCount > 4),超过则标记为 ERROR。失败任务由定时任务补偿。
public interface ITradePort {
String groupBuyNotify(NotifyTaskEntity notifyTask);
}
domain 只定义接口,infrastructure 实现(HTTP 回调或 MQ 发送)。这是 DDD 端口-适配器架构的体现。
public interface ITagService {
void execTagBatchJob(String tagId, String batchId);
}
@Service
public class TagService implements ITagService {
@Override
public void execTagBatchJob(String tagId, String batchId) {
// 1. 查询批次任务配置
CrowdTagsJobEntity job = repository.queryCrowdTagsJobEntity(tagId, batchId);
// 2. 采集用户数据(实际生产中由数仓团队通过脚本写入)
// 3. 模拟写入人群标签
List<String> userIdList = Arrays.asList("xiaofuge", "liergou", ...);
for (String userId : userIdList) {
repository.addCrowdTagsUserId(tagId, userId);
}
// 4. 更新标签统计
repository.updateCrowdTagsStatistics(tagId, userIdList.size());
}
}
说明:当前版本的人群数据是硬编码的测试数据。真实生产环境中,execTagBatchJob 会从大数据平台拉取符合条件的用户列表,批量写入 crowd_tags_detail 表。
@Data @Builder
public class CrowdTagsJobEntity {
private Integer tagType; // 标签类型(参与量/消费金额)
private String tagRule; // 标签规则(如 "N>3" 表示参与3次以上)
private Date statStartTime; // 统计开始时间
private Date statEndTime; // 统计结束时间
}
| 模式 | 应用位置 | 类 |
|---|---|---|
| 策略模式 | 折扣计算 | IDiscountCalculateService + ZJ/MJ/ZK/NCalculateService |
| 策略模式 | 退单处理 | IRefundOrderStrategy + 3 个策略类 |
| 模板方法 | 折扣计算 | AbstractDiscountCalculateService → doCalculate() |
| 模板方法 | 退单 | AbstractRefundOrderStrategy → refundOrder() / reverseStock() |
| 模板方法 | 决策树框架 | AbstractMultiThreadStrategyRouter → doApply() / get() |
| 决策树 | 试算链 | Root → Switch → Market → Tag → End/Error |
| 责任链 | 锁单过滤 | ActivityUsability → UserTakeLimit → TeamStockOccupy |
| 责任链 | 结算过滤 | SCRule → OutTradeNo → Settable → EndRule |
| 责任链 | 退单过滤 | DataNode → UniqueRefund → RefundOrder |
| 聚合模式 | 实体打包 | GroupBuyOrderAggregate / GroupBuyRefundAggregate / GroupBuyTeamSettlementAggregate |
| 工厂模式 | 链创建 | TradeLockRuleFilterFactory / TradeSettlementRuleFilterFactory / TradeRefundRuleFilterFactory |
| 端口-适配器 | 仓储/外部接口 | IActivityRepository / ITradeRepository / ITagRepository / ITradePort |
| 动态上下文 | 状态传递 | DefaultActivityStrategyFactory.DynamicContext / 各 Factory 的 DynamicContext |
| 多线程异步 | 试算数据加载 | MarketNode#multiThread() + FutureTask / CompletableFuture |
| 枚举多态 | 退单策略分发 | RefundTypeEnumVO 的抽象 matches() 方法 |
将整个系统的数据流串联起来,一个完整的"用户从浏览到成功拼团"流程如下:
┌─ 用户打开商品详情页 ──────────────────────────────────────────────┐
│ │
│ 1. trigger → IIndexGroupBuyMarketService.indexMarketTrial() │
│ ├── RootNode: 参数校验 │
│ ├── SwitchNode: 降级/切量判断 │
│ ├── MarketNode: 多线程加载活动配置+商品信息 │
│ │ → discountCalculateServiceMap.get(marketPlan).calculate() │
│ │ → 写入 deductionPrice, payPrice 到 DynamicContext │
│ ├── TagNode: 人群标签过滤 (visible/enable) │
│ └── EndNode: 组装 TrialBalanceEntity 返回 │
│ → 前端展示「拼团价 XX 元,N人成团」 │
│ │
├─ 用户点击「参与拼团」下单 ────────────────────────────────────────│
│ │
│ 2. trigger → ITradeLockOrderService.lockMarketPayOrder() │
│ ├── ActivityUsabilityRuleFilter: 活动是否生效 │
│ ├── UserTakeLimitRuleFilter: 用户是否达到参与上限 │
│ ├── TeamStockOccupyRuleFilter: Redis 抢占库存 │
│ └── repository.lockMarketPayOrder(aggregate) │
│ → 写入 group_buy_order + group_buy_order_list │
│ │
├─ 用户完成支付 ───────────────────────────────────────────────────│
│ │
│ 3. trigger → ITradeSettlementOrderService.settlementMarketPayOrder│
│ ├── SCRuleFilter: 渠道黑名单 │
│ ├── OutTradeNoRuleFilter: 外部单号存在/未退单 │
│ ├── SettableRuleFilter: 支付时间在有效期内 │
│ ├── EndRuleFilter: 组装结算数据 │
│ ├── repository.settlementMarketPayOrder(aggregate) │
│ │ → 更新 group_buy_order_list (lockCount+1, completeCount+1) │
│ │ → 判断是否达成目标人数 → 如达成,更新 group_buy_order 状态 │
│ │ → 如达成,创建 notify_task 回调记录 │
│ └── 异步: tradeTaskService.execNotifyJob(notifyTask) │
│ → HTTP 回调通知外部系统「拼团成功」 │
│ │
├─ 用户退单 ────────────────────────────────────────────────────────│
│ │
│ 4. trigger → ITradeRefundOrderService.refundOrder() │
│ ├── DataNodeFilter: 查询订单+团队数据 │
│ ├── UniqueRefundNodeFilter: 幂等检查 │
│ ├── RefundOrderNodeFilter: 状态组合 → 策略分发 │
│ │ ├── UNPAID_UNLOCK → Unpaid2RefundStrategy │
│ │ ├── PAID_UNFORMED → Paid2RefundStrategy │
│ │ └── PAID_FORMED → PaidTeam2RefundStrategy │
│ └── 异步: MQ → restoreTeamLockStock → Redis 库存恢复 │
│ │
└────────────────────────────────────────────────────────────────────┘
依赖倒置:domain 定义 IActivityRepository、ITradeRepository、ITradePort 接口,infrastructure 实现它们。domain 不依赖具体技术。
聚合根:GroupBuyOrderAggregate 不是简单的"把几个实体放一起",它代表了"锁单"这个业务操作的事务边界——infrastructure 通过它一次写入多张表。
值对象 vs 实体:
GroupBuyTeamEntity.teamId、MarketPayOrderEntity.orderIdGroupBuyProgressVO、NotifyConfigVO端口-适配器:ITradePort 是端口(domain 定义),HTTP 实现是适配器(infrastructure 提供)。
| 特性 | 决策树(xfg-wrench tree) | 责任链(xfg-wrench link) |
|---|---|---|
| 路由方式 | 每个节点的 get() 决定下一个节点 | 固定链式顺序 |
| 中断方式 | 路由到 ErrorNode | 抛异常中断 |
| 适用场景 | 有多分支条件的流程 | 有固定顺序的规则校验 |
| 代表 | 试算流程(Root→Switch→Market→Tag→End) | 锁单/结算/退单 |
在每个 Filter 链/决策树中都有一个 DynamicContext 内部类,它的作用是:
new DynamicContext(),不存在并发问题@Resource
private Map<String, IDiscountCalculateService> discountCalculateServiceMap;
Spring 会自动把所有 IDiscountCalculateService 的实现 bean 注入到这个 Map 中,key 是 bean 名称(@Service("ZJ") → key="ZJ"),value 是实例。这是 Spring 提供的策略模式开箱即用的能力。
RefundTypeEnumVO 展示了 Java 枚举的高级用法——每个枚举值可以覆写抽象方法,实现不同的行为。配合 strategy 字段指向 Spring bean 名称,实现了从数据状态到策略的一步映射。
如果你想在这个项目里新增一种拼团玩法,最可能的改动路径是:
DiscountTypeEnum 或 marketPlan 里新增类型XxxCalculateService extends AbstractDiscountCalculateService这就是 DDD + 策略/责任链带来的开闭原则——新增功能只需加新类,无需修改已有代码。
下一篇建议阅读:
group-buy-market-infrastructure模块 —— 了解 domain 层定义的接口如何被具体的数据库/Redis/MQ 实现。
知识笔记会随着实践和认知变化持续更新,不代表最终结论。