Knowledge note
正在加载知识笔记
正在加载知识笔记
Knowledge note
模块定位:DDD 基础设施层 — domain 接口的技术实现
包路径:cn.bugstack.infrastructure
依赖关系:依赖domain(实现其接口)+types+ Spring Boot + MyBatis + Redisson + RabbitMQ + OkHttp
源码文件:约 35 个 Java 类 + 10 个 MyBatis XML
如果说 domain 模块定义了"应该做什么",那 infrastructure 模块就是"具体怎么做"——它是 domain 层所有接口的技术实现。
domain 层定义接口 infrastructure 层实现
┌──────────────────────┐ ┌──────────────────────────┐
│ IActivityRepository │ ← implements ──│ ActivityRepository │
│ ITradeRepository │ ← implements ──│ TradeRepository │
│ ITagRepository │ ← implements ──│ TagRepository │
│ ITradePort │ ← implements ──│ TradePort │
└──────────────────────┘ └──────────────────────────┘
│
使用 DAO / Redis / MQ / HTTP
infrastructure 模块只做技术实现,不包含业务逻辑:
group-buy-market-infrastructure/
└── src/main/java/cn/bugstack/infrastructure/
│
├── adapter/
│ ├── repository/ # ★ Repository 实现(核心)
│ │ ├── AbstractRepository.java # 通用缓存模板基类
│ │ ├── ActivityRepository.java # 活动仓储(含缓存策略)
│ │ ├── TradeRepository.java # 交易仓储(最复杂:锁单/结算/退单)
│ │ └── TagRepository.java # 标签仓储(DB + Redis BitSet)
│ └── port/
│ └── TradePort.java # ★ 回调通知端口(HTTP/MQ 分发 + 分布式锁)
│
├── dao/ # ★ MyBatis DAO 接口
│ ├── IGroupBuyActivityDao.java # 拼团活动表
│ ├── IGroupBuyDiscountDao.java # 折扣配置表
│ ├── IGroupBuyOrderDao.java # 拼团订单(组队)表
│ ├── IGroupBuyOrderListDao.java # 拼团订单明细表
│ ├── INotifyTaskDao.java # 通知任务表
│ ├── ISkuDao.java # 商品表
│ ├── ISCSkuActivityDao.java # 渠道商品关联表
│ ├── ICrowdTagsDao.java # 人群标签表
│ ├── ICrowdTagsDetailDao.java # 人群标签明细表
│ ├── ICrowdTagsJobDao.java # 人群标签任务表
│ └── po/ # 持久化对象(与DB表一一对应)
│ ├── GroupBuyActivity.java
│ ├── GroupBuyDiscount.java
│ ├── GroupBuyOrder.java
│ ├── GroupBuyOrderList.java
│ ├── NotifyTask.java
│ ├── Sku.java
│ ├── SCSkuActivity.java
│ ├── CrowdTags.java
│ ├── CrowdTagsDetail.java
│ ├── CrowdTagsJob.java
│ └── base/Page.java
│
├── redis/ # Redis 服务层
│ ├── IRedisService.java # Redis 操作接口(35+ 方法)
│ └── RedissonService.java # Redisson 实现
│
├── dcc/
│ └── DCCService.java # ★ 动态配置中心(降级/切量/黑名单/缓存)
│
├── event/
│ └── EventPublisher.java # RabbitMQ 消息发布
│
└── gateway/
└── GroupBuyNotifyService.java # HTTP 回调服务(OkHttp)
配合 MyBatis XML(位于 app 模块 resources):
group-buy-market-app/src/main/resources/mybatis/mapper/
├── group_buy_activity_mapper.xml
├── group_buy_discount_mapper.xml
├── group_buy_order_mapper.xml
├── group_buy_order_list_mapper.xml
├── nofify_task_mapper.xml ← 注意拼写
├── sku_mapper.xml
├── sc_sku_activity_mapper.xml
├── crowd_tags_mapper.xml
├── crowd_tags_detail_mapper.xml
└── crowd_tags_job_mapper.xml
infrastructure 内部的数据流转是严格的四层单向调用:
MySQL 数据库表
↕ (MyBatis XML 映射)
PO (Persistent Object) ← 与 DB 表字段一一对应
↕
DAO (@Mapper 接口) ← MyBatis 自动生成代理实现
↕
Repository (@Repository) ← 实现 domain 的 IXXXRepository 接口
↕ (负责 PO → Entity 转换 + 事务管理)
Domain Entity / VO ← domain 层定义的业务模型
每一层的作用:
| 层 | 作用 | 关键注解 |
|---|---|---|
| PO | 数据库行在 Java 中的映射,字段=列 | @Data @Builder |
| DAO | 声明 SQL 方法,不写实现 | @Mapper |
| Repository | 实现 domain 接口,调 DAO,转换类型,管理事务 | @Repository @Transactional |
| Entity | 属于 domain 层,infrastructure 只构造它 | @Data @Builder(在 domain 定义) |
// PO(infrastructure 层)—— 与数据库字段完全相同,可以用 @Table 等注解
public class GroupBuyOrder {
private Long id; // 自增主键
private String teamId;
private Integer status; // 0/1/2 存 int
private Date createTime; // 有创建时间
private Date updateTime; // 有更新时间
}
// Entity(domain 层)—— 业务语义更强,没有DB才有的字段
public class GroupBuyTeamEntity {
private String teamId;
private GroupBuyOrderEnumVO status; // 枚举,不是 int
// 没有 id, createTime, updateTime
}
为什么分开?
三个 Repository 实现了 domain 层定义的三个仓储接口:
| Repository | 实现的 Domain 接口 | 主要操作的表 |
|---|---|---|
ActivityRepository | IActivityRepository | group_buy_activity, group_buy_discount, sku, sc_sku_activity, group_buy_order, group_buy_order_list |
TradeRepository | ITradeRepository | group_buy_order, group_buy_order_list, notify_task |
TagRepository | ITagRepository | crowd_tags, crowd_tags_detail, crowd_tags_job |
这是 infrastructure 模块最有复用价值的代码。
public abstract class AbstractRepository {
@Resource
protected IRedisService redisService;
@Resource
protected DCCService dccService;
// 核心方法:缓存优先读取
protected <T> T getFromCacheOrDb(String cacheKey, Supplier<T> dbFallback) {
// 1. 检查缓存开关
if (dccService.isCacheOpenSwitch()) {
// 2. 先从 Redis 获取
T cacheResult = redisService.getValue(cacheKey);
if (null != cacheResult) {
return cacheResult; // 缓存命中 → 直接返回
}
// 3. 缓存未命中 → 查数据库
T dbResult = dbFallback.get();
if (null == dbResult) return null;
// 4. 写入缓存
redisService.setValue(cacheKey, dbResult);
return dbResult;
} else {
// 缓存降级 → 直接走数据库
return dbFallback.get();
}
}
}
1. Supplier 函数式接口做回调
getFromCacheOrDb(cacheKey, () -> groupBuyActivityDao.queryValidGroupBuyActivityId(activityId));
// ↑ 只有缓存未命中时才执行这个 Lambda
2. 缓存开关降级
通过 dccService.isCacheOpenSwitch() 可以在不重启服务的情况下关闭缓存——如果 Redis 出问题,运维可以直接关掉缓存保证服务可用。
3. ActivityRepository 使用示例
// 活动信息缓存
GroupBuyActivity groupBuyActivityRes = getFromCacheOrDb(
GroupBuyActivity.cacheRedisKey(activityId), // key: 含全限定类名
() -> groupBuyActivityDao.queryValidGroupBuyActivityId(activityId)
);
// 折扣信息缓存(不同的缓存 key)
GroupBuyDiscount groupBuyDiscountRes = getFromCacheOrDb(
GroupBuyDiscount.cacheRedisKey(discountId),
() -> groupBuyDiscountDao.queryGroupBuyActivityDiscountByDiscountId(discountId)
);
// GroupBuyActivity.java (PO)
public static String cacheRedisKey(Long activityId) {
return "group_buy_market_cn.bugstack.infrastructure.dao.po.GroupBuyActivity_" + activityId;
}
// GroupBuyDiscount.java (PO)
public static String cacheRedisKey(String discountId) {
return "group_buy_market_cn.bugstack.infrastructure.dao.po.GroupBuyDiscount_" + discountId;
}
Key 包含全限定类名,确保不同 PO 的缓存 key 不会冲突。实际生产中可以简化为 group_buy_market:activity:{id}。
这是 infrastructure 模块最复杂、代码量最大的类(~700 行),实现了 ITradeRepository 的全部 20+ 个方法。
| 方法 | 事务 | 操作的表 | 业务场景 |
|---|---|---|---|
queryMarketPayOrderEntityByOutTradeNo | 无 | group_buy_order_list | 查询订单详情 |
lockMarketPayOrder | 有 (500ms) | group_buy_order + group_buy_order_list | 锁单(首次插入 / 后续更新) |
queryGroupBuyProgress | 无 | group_buy_order | 查询成团进度 |
queryGroupBuyActivityEntityByActivityId | 无 | group_buy_activity | 查询活动信息 |
queryOrderCountByActivityId | 无 | group_buy_order_list | 用户已参与次数 |
queryGroupBuyTeamByTeamId | 无 | group_buy_order | 查询队伍详情 |
settlementMarketPayOrder | 有 (5s) | group_buy_order_list + group_buy_order + notify_task | 支付结算 |
occupyTeamStock | 无 | 仅 Redis | 组队库存抢占 |
recoveryTeamStock | 无 | 仅 Redis | 库存恢复 |
unpaid2Refund | 有 (5s) | group_buy_order_list + group_buy_order + notify_task | 未支付退单 |
paid2Refund | 有 (5s) | 同上 | 已支付退单 |
paidTeam2Refund | 有 (5s) | 同上 | 已成团退单 |
refund2AddRecovery | 无 | 仅 Redis | 退单恢复库存 |
queryTimeoutUnpaidOrderList | 无 | group_buy_order_list + group_buy_order | 超时未支付订单 |
@Transactional(timeout = 500)
public MarketPayOrderEntity lockMarketPayOrder(GroupBuyOrderAggregate aggregate) {
// 1. teamId 为空 → 首次开团,INSERT group_buy_order
if (StringUtils.isBlank(teamId)) {
teamId = RandomStringUtils.randomNumeric(8); // 生成8位队伍ID
groupBuyOrderDao.insert(groupBuyOrder); // lockCount = 1, completeCount = 0
} else {
// 2. teamId 非空 → 加入已有团,UPDATE lockCount + 1
int updateAddTargetCount = groupBuyOrderDao.updateAddLockCount(teamId);
if (1 != updateAddTargetCount) {
throw new AppException(ResponseCode.E0005); // 更新失败 → 队伍已满
}
}
// 3. 插入订单明细(每个用户一条)
GroupBuyOrderList groupBuyOrderListReq = GroupBuyOrderList.builder()
.bizId(activityId + "_" + userId + "_" + (userTakeOrderCount + 1)) // 唯一业务ID
...build();
groupBuyOrderListDao.insert(groupBuyOrderListReq); // DuplicateKeyException → E0003
// 4. 返回 MarketPayOrderEntity
return MarketPayOrderEntity.builder()...build();
}
关键设计点:
bizId 唯一索引:activityId_userId_参与次数 作为数据库唯一索引,确保一个用户在同一活动上只能参与 takeLimitCount 次。依赖 DuplicateKeyException 做并发控制。
teamId 生成:RandomStringUtils.randomNumeric(8) 生成 8 位数字队伍ID(生产环境应用雪花算法)。
@Transactional(timeout = 5000)
public NotifyTaskEntity settlementMarketPayOrder(GroupBuyTeamSettlementAggregate aggregate) {
// 1. 更新订单明细状态:CREATE → COMPLETE
int updateOrderListStatusCount = groupBuyOrderListDao.updateOrderStatus2COMPLETE(req);
// 受影响行数 ≠ 1 → 抛异常
// 2. 更新组队完成数量:completeCount + 1
int updateAddCount = groupBuyOrderDao.updateAddCompleteCount(teamId);
// 3. 判断是否最后一个(targetCount - completeCount == 1)
if (groupBuyTeamEntity.getTargetCount() - groupBuyTeamEntity.getCompleteCount() == 1) {
// 更新组队状态为 COMPLETE
groupBuyOrderDao.updateOrderStatus2COMPLETE(teamId);
// 查询外部单号列表
List<String> outTradeNoList = groupBuyOrderListDao
.queryGroupBuyCompleteOrderOutTradeNoListByTeamId(teamId);
// 写入 notify_task 表(拼团成功通知)
NotifyTask notifyTask = buildNotifyTask(...);
notifyTaskDao.insert(notifyTask);
}
}
注释中的并发问题思考(源码第 262 行):
// 【面试题,这个地方可能会有一个并发情况,
// 就是多个用户拿到的 groupBuyTeamEntity.getCompleteCount() 是同一个值怎么办?】
// 【方式1;可以给调用 settlementMarketPayOrder 结算方法的地方,添加一个分布式锁】
// 【方式2;这部分结算,只做数据库的更新操作,以及发送mq,之后在消费mq的地方,做结算】
// 【方式3;增加一个定时job任务补偿,检索订单量够,但没有结算的拼团组队记录】
public boolean occupyTeamStock(String teamStockKey, String recoveryTeamStockKey,
Integer target, Integer validTime) {
// 1. 获取失败恢复量
Long recoveryCount = redisService.getAtomicLong(recoveryTeamStockKey);
// 2. incr + 1(因为已有1个占用量)
long occupy = redisService.incr(teamStockKey) + 1;
// 3. 超出库存限制 → 回退 decr
if (occupy > target + recoveryCount) {
redisService.decr(teamStockKey);
return false;
}
// 4. 加 Redis 锁做兜底(防止集群问题导致 incr 值重复)
String lockKey = teamStockKey + "_" + occupy;
Boolean lock = redisService.setNx(lockKey, validTime + 60, TimeUnit.MINUTES);
return lock;
}
public GroupBuyActivityDiscountVO queryGroupBuyActivityDiscountVO(Long activityId) {
// 两级缓存:先查 Redis,miss 再查 DB
GroupBuyActivity groupBuyActivityRes = getFromCacheOrDb(
GroupBuyActivity.cacheRedisKey(activityId),
() -> groupBuyActivityDao.queryValidGroupBuyActivityId(activityId));
GroupBuyDiscount groupBuyDiscountRes = getFromCacheOrDb(
GroupBuyDiscount.cacheRedisKey(discountId),
() -> groupBuyDiscountDao.queryGroupBuyActivityDiscountByDiscountId(discountId));
// PO → VO 转换
return GroupBuyActivityDiscountVO.builder()...build();
}
两级缓存:活动信息和折扣信息分别缓存,因为它们对应不同的表(group_buy_activity 和 group_buy_discount),各自的失效策略不同。
public boolean isTagCrowdRange(String tagId, String userId) {
RBitSet bitSet = redisService.getBitSet(tagId);
if (!bitSet.isExists()) return true; // 不存在 → 无限制
return bitSet.get(redisService.getIndexFromUserId(userId));
}
为什么用 BitSet?
| 方案 | 内存占用(1000万用户) | 查询速度 |
|---|---|---|
| Redis Set | ~800MB | O(1) |
| Redis BitSet | ~1.2MB | O(1) |
BitSet 把 userId 的 MD5 哈希映射到一个 bit 位,内存占用极小,适合大批量人群标签场景。
public List<UserGroupBuyOrderDetailEntity> queryInProgress...ListByRandom(...) {
// 1. 从 DB 查 2 倍数量的数据
List<GroupBuyOrderList> orderLists = dao.queryByRandom(req); // count = randomCount * 2
// 2. 如果查到的数据多于需要的 → 随机打乱 + 截取
if (orderLists.size() > randomCount) {
Collections.shuffle(orderLists);
orderLists = orderLists.subList(0, randomCount);
}
// 3. 批量查组队信息 → 组装返回
...
}
这是典型的前端"随机展示"实现——先取 2 倍数据,再 shuffle + 截取,既保证了随机性,又避免每次查全表。
@Service
public class TradePort implements ITradePort {
@Override
public String groupBuyNotify(NotifyTaskEntity notifyTask) throws Exception {
// 1. 分布式锁 — 防止多台机器重复执行
RLock lock = redisService.getLock(notifyTask.lockKey());
if (lock.tryLock(3, 0, TimeUnit.SECONDS)) {
try {
// 2. HTTP 回调
if (NotifyTypeEnumVO.HTTP.getCode().equals(notifyTask.getNotifyType())) {
if (StringUtils.isBlank(notifyTask.getNotifyUrl()) || "暂无".equals(notifyTask.getNotifyUrl())) {
return NotifyTaskHTTPEnumVO.SUCCESS.getCode(); // 无效URL直接成功
}
groupBuyNotifyService.groupBuyNotify(notifyTask.getNotifyUrl(), notifyTask.getParameterJson());
return NotifyTaskHTTPEnumVO.SUCCESS.getCode();
}
// 3. MQ 回调
if (NotifyTypeEnumVO.MQ.getCode().equals(notifyTask.getNotifyType())) {
publisher.publish(notifyTask.getNotifyMQ(), notifyTask.getParameterJson());
return NotifyTaskHTTPEnumVO.SUCCESS.getCode();
}
} finally {
lock.unlock();
}
}
return NotifyTaskHTTPEnumVO.NULL.getCode(); // 没抢到锁 → 空执行
}
}
分布式锁的必要性:lockKey = "notify_job_lock_key_" + uuid。定时任务 GroupBuyNotifyJob 可能在多台机器上同时触发,分布式锁确保同一通知只被处理一次。
超时打断:lock.tryLock(3, 0, TimeUnit.SECONDS) 3 秒超时 + 0 等待时间,拿不到锁立即返回 NULL。
"暂无"的特殊处理:notifyUrl 为 "暂无" 直接返回成功——这是运营配置的占位符语义,不是真需要回调。
domain / infrastructure
│
▼
IRedisService ← 自定义接口(35+ 方法)
│
RedissonService ← 实现,委托 RedissonClient
│
RedissonClient ← Redisson 框架
│
Redis Server
直接使用 RedissonClient 的问题:
自定义 IRedisService 的好处:
setValue / getValue 里统一加日志、监控RedissonService,不影响业务代码| 类别 | 方法 | 底层 Redisson API |
|---|---|---|
| KV 存储 | setValue, getValue, remove, isExists | RBucket |
| 计数器 | incr, decr, setAtomicLong, getAtomicLong | RAtomicLong |
| 分布式锁 | getLock, getFairLock, getReadWriteLock | RLock |
| 集合 | addToSet, isSetMember | RSet |
| 列表 | addToList, getFromList | RList |
| 哈希 | addToMap, getFromMap, getMap | RMap |
| 队列 | getQueue, getBlockingQueue, getDelayedQueue | RQueue |
| BitSet | getBitSet | RBitSet |
| 信号量 | getSemaphore, getCountDownLatch | RSemaphore |
| 布隆过滤器 | getBloomFilter | RBloomFilter |
default int getIndexFromUserId(String userId) {
MessageDigest md = MessageDigest.getInstance("MD5");
byte[] hashBytes = md.digest(userId.getBytes(StandardCharsets.UTF_8));
BigInteger bigInt = new BigInteger(1, hashBytes); // 正整数
return bigInt.mod(BigInteger.valueOf(Integer.MAX_VALUE)) // 取模到 0~2^31-1
.intValue();
}
MD5 哈希 + 取模 → 把任意 userId 映射到 0~21亿 的整数范围,作为 BitSet 的 bit 位索引。
@Service
public class DCCService {
@DCCValue("downgradeSwitch:0") // 降级开关,默认 0(关闭)
private String downgradeSwitch;
@DCCValue("cutRange:100") // 切量范围,默认 100(全量)
private String cutRange;
@DCCValue("scBlacklist:s02c02") // 渠道黑名单,默认 "s02c02"
private String scBlacklist;
@DCCValue("cacheSwitch:0") // 缓存开关,默认 0(开启)
private String cacheOpenSwitch;
}
@DCCValue 注解来自 xfg-wrench 框架,类似 Apollo/Nacos 的 @Value 注解,但值存储在 Redis 中。格式:@DCCValue("key:defaultValue")。
好处:不重启服务就能改配置,Redis 值更新后自动注入字段。
| 配置项 | 默认值 | 作用 | 使用者 |
|---|---|---|---|
downgradeSwitch | "0" | "1" = 降级(抛出E0003) | SwitchNode |
cutRange | "100" | 全量百分比(用户hash%100 <= 此值则可见) | SwitchNode |
scBlacklist | "s02c02" | 逗号分隔的渠道黑名单 | SCRuleFilter |
cacheOpenSwitch | "0" | "0" = 开启缓存,"1" = 关闭(降级直查DB) | AbstractRepository |
public boolean isCutRange(String userId) {
int hashCode = Math.abs(userId.hashCode());
int lastTwoDigits = hashCode % 100; // 取 hash 的后两位 0~99
return lastTwoDigits <= Integer.parseInt(cutRange); // cutRange=30 → 30% 的用户可见
}
这是典型的灰度发布方案:通过 hash 取模将用户分桶,逐步放量。cutRange=10 只有 10% 用户可见拼团活动,cutRange=100 全量。
@Component
public class EventPublisher {
@Autowired
private RabbitTemplate rabbitTemplate;
@Value("${spring.rabbitmq.config.producer.exchange}")
private String exchangeName; // "group_buy_market_exchange"
public void publish(String routingKey, String message) {
rabbitTemplate.convertAndSend(exchangeName, routingKey, message, m -> {
// 设置消息持久化(服务重启不丢失)
m.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return m;
});
}
}
两种 MQ 消息路由:
| Routing Key | 消息内容 | 消费者(trigger 模块) |
|---|---|---|
topic.team_success | 拼团成功通知(teamId + outTradeNoList) | TeamSuccessTopicListener |
topic.team_refund | 退单通知(type + userId + orderId + teamId) | RefundSuccessTopicListener |
@Service
public class GroupBuyNotifyService {
@Resource
private OkHttpClient okHttpClient; // 由 app 模块的 OKHttpClientConfig 创建
public String groupBuyNotify(String apiUrl, String notifyRequestDTOJSON) {
// 1. 构造 JSON 请求体
MediaType mediaType = MediaType.parse("application/json");
RequestBody body = RequestBody.create(mediaType, notifyRequestDTOJSON);
// 2. POST 请求
Request request = new Request.Builder()
.url(apiUrl).post(body)
.addHeader("content-type", "application/json")
.build();
// 3. 同步执行
Response response = okHttpClient.newCall(request).execute();
return response.body().string();
}
}
拼团成团后,系统主动 POST 回调通知外部业务系统(如订单系统),告知拼团结果。调用方在 LockMarketPayOrderRequestDTO.setNotifyUrl() 中指定回调地址。
| DAO 接口 | 对应表 | 方法数 | 关键方法 |
|---|---|---|---|
IGroupBuyActivityDao | group_buy_activity | 4 | queryValidGroupBuyActivityId |
IGroupBuyDiscountDao | group_buy_discount | 2 | queryGroupBuyActivityDiscountByDiscountId |
IGroupBuyOrderDao | group_buy_order | 14 | insert, updateAddLockCount, updateAddCompleteCount, unpaid2Refund, paid2Refund, paidTeam2Refund |
IGroupBuyOrderListDao | group_buy_order_list | 11 | insert, updateOrderStatus2COMPLETE, queryInProgress...ByUserId, queryTimeoutUnpaidOrderList |
INotifyTaskDao | notify_task | 6 | insert, queryUnExecutedNotifyTaskList, updateNotifyTaskStatusSuccess |
ISkuDao | sku | 1 | querySkuByGoodsId |
ISCSkuActivityDao | sc_sku_activity | 1 | querySCSkuActivityBySCGoodsId |
ICrowdTagsDao | crowd_tags | 1 | updateCrowdTagsStatistics |
ICrowdTagsDetailDao | crowd_tags_detail | 1 | addCrowdTagsUserId |
ICrowdTagsJobDao | crowd_tags_job | 1 | queryCrowdTagsJob |
1. PO 自带缓存 Key 方法
// GroupBuyActivity.java
public static String cacheRedisKey(Long activityId) {
return "group_buy_market_cn.bugstack.infrastructure.dao.po.GroupBuyActivity_" + activityId;
}
把缓存 key 的生成逻辑放在 PO 类内部,避免了 Repository 里写硬编码字符串。
2. GroupBuyOrderList 继承 Page
public class GroupBuyOrderList extends Page {
// ...字段
}
// Page.java
public class Page {
private Integer count; // 用于 limit 的通用分页参数
}
查询时设置 count 字段 → MyBatis XML 中用 LIMIT #{count} 实现分页。
XML 文件位于 app 模块的 resources/mybatis/mapper/ 下(因为 application.yml 中有 mybatis.mapper-locations: classpath:/mybatis/mapper/*.xml)。
以 group_buy_order_mapper.xml 为例:
<mapper namespace="cn.bugstack.infrastructure.dao.IGroupBuyOrderDao">
<insert id="insert" parameterType="cn.bugstack.infrastructure.dao.po.GroupBuyOrder">
INSERT INTO group_buy_order (team_id, activity_id, source, channel, ...)
VALUES (#{teamId}, #{activityId}, #{source}, #{channel}, ...)
</insert>
</mapper>
<update id="updateAddLockCount" parameterType="java.lang.String">
UPDATE group_buy_order
SET lock_count = lock_count + 1, update_time = now()
WHERE team_id = #{teamId}
AND lock_count < target_count
</update>
lock_count < target_count 是 CAS 条件——如果拼团已满(lock_count = target_count),affected rows = 0,上层代码通过 if (1 != updateCount) 抛出异常。
<select id="queryGroupBuyProgressByTeamIds" resultType="...">
SELECT * FROM group_buy_order
WHERE team_id IN
<foreach item="teamId" collection="teamIds" open="(" separator="," close=")">
#{teamId}
</foreach>
</select>
用于把多个 teamId 一次性查出,避免 N+1 查询。
以锁单操作为例,展示 infrastructure 如何贯穿各层:
1. trigger 层调用 → app 层调用 → domain LockMarketPayOrder
2. domain 层 TradeLockOrderService.lockMarketPayOrder()
→ 责任链校验通过
→ 调用 repository.lockMarketPayOrder(aggregate) ← 这是 domain 的 ITradeRepository 接口
3. infrastructure 层 TradeRepository.lockMarketPayOrder()
│
├── 3a. teamId 为空 → groupBuyOrderDao.insert(po) → INSERT INTO group_buy_order
│ teamId 非空 → groupBuyOrderDao.updateAddLockCount(teamId)
│ → UPDATE group_buy_order SET lock_count = lock_count + 1
│
├── 3b. groupBuyOrderListDao.insert(po) → INSERT INTO group_buy_order_list
│ (id, userId, teamId, orderId, activityId, ..., bizId, ...)
│ 如果 bizId 重复 → DuplicateKeyException → AppException(E0003)
│
└── 3c. 返回 MarketPayOrderEntity(orderId, prices, teamId)
→ domain 层拿到 Entity → app 层转成 ResponseDTO → trigger 返回 HTTP 响应
| 模式/实践 | 位置 | 说明 |
|---|---|---|
| 依赖倒置 | TradeRepository implements ITradeRepository | infrastructure 实现 domain 的接口 |
| 模板方法 | AbstractRepository.getFromCacheOrDb() | 定义缓存读取骨架,子类只需传 key 和查询函数 |
| 回调函数(Supplier) | getFromCacheOrDb(key, () -> dao.query()) | Lambda 延迟执行,缓存命中不查 DB |
| 兜底降级 | dccService.isCacheOpenSwitch() | 缓存不可用时自动走 DB |
| CAS 更新 | UPDATE ... WHERE lock_count < target_count | SQL 层面的乐观锁 |
| 唯一索引防重 | bizId 字段 | 数据库层面防止重复参与 |
| 分布式锁 | redisService.getLock(notifyTask.lockKey()) | 防止多机重复回调 |
| BitSet 人群过滤 | redisService.getBitSet(tagId) | 极低内存的会员过滤 |
| Shuffle 随机 | Collections.shuffle() + subList() | 前端随机展示 |
很多初学 Spring Boot 的开发者会把 DAO 和 Repository 混为一谈:
insert, queryByTeamIdlockMarketPayOrder, settlementMarketPayOrderRepository 调用 DAO,但比 DAO 更"业务化"——一个 Repository 方法可能操作多张表、包含事务、做 PO→Entity 转换。
如果 PO 和 Entity 是同一个类,会出现:
createTime、updateTime(只对 DB 有意义)分开后,换数据库只影响 infrastructure 层的 PO + DAO + XML,domain 完全不受影响。
AbstractRepository.getFromCacheOrDb() 处理了缓存穿透(大量请求同时查一个不存在的 key)的情况——通过 DCCService 的缓存开关做降级。但它没有处理缓存击穿(热点 key 过期瞬间大量请求打 DB)——生产环境可以加分布式锁或设置随机过期时间。
lock.tryLock(3, 0, TimeUnit.SECONDS); // 3秒超时,0等待
try {
// 业务逻辑
} finally {
if (lock.isLocked() && lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
关键点:
tryLock 不是 lock——拿不到锁立即返回,不会阻塞finally 中先判断 isHeldByCurrentThread() 再 unlock——避免释放别人的锁notifyTask.lockKey())——颗粒度精确到单条记录UPDATE group_buy_order
SET lock_count = lock_count + 1
WHERE team_id = #{teamId}
AND lock_count < target_count -- ← 这就是乐观锁!
不需要 SELECT FOR UPDATE,不需要 Redis 锁——直接用 SQL 条件判断。affected rows = 0 意味着库存已满,上层抛异常即可。
下一篇建议阅读:
group-buy-market-trigger模块 —— 了解 HTTP Controller、定时任务、MQ Listener 如何成为系统的外部入口。
知识笔记会随着实践和认知变化持续更新,不代表最终结论。