家纺AI生图平台 · 积分与会员系统技术设计¶
文档定位:目标架构(事件溯源 + 快照 + 幂等 + 对账)。现行实现为「会员方案 → 单一
usage_quotas扣点」,见 membership-agent-sales-guide.md;本文描述金融级账本目标,不代表已全部落地。目标:确保账目 准确性 · 持久性 · 健壮性 · 可回溯 · 幂等 · 日查,达到金融级可靠性
目录¶
一、设计原则¶
┌──────────────────────────────────────────────────┐
│ 金融级账目系统五原则 │
├──────────────────────────────────────────────────┤
│ 1. 强一致性 —— 任何时刻余额=SUM(所有流水) │
│ 2. 幂等性 —— 同一操作重复执行,结果不变 │
│ 3. 可追溯性 —— 每一分积分的来龙去脉完整记录 │
│ 4. 不可篡改 —— 已确认的流水不可修改,只能冲正 │
│ 5. 可对账 —— 任意时间点可自动核算余额与流水一致性 │
└──────────────────────────────────────────────────┘
核心架构选择:事件溯源(Event Sourcing)+ 快照(Snapshot)
- 所有积分变动以「事件」形式写入不可变事件表
- 当前余额由事件重放(Replay)计算得出,并定期生成快照加速查询
- 余额快照表与事件表相互校验,任何时刻可查差异
二、数据模型设计¶
2.1 核心表结构¶
(1)积分事件表(不可变,只追加)points_event¶
CREATE TABLE points_event (
event_id BIGINT PRIMARY KEY, -- 全局唯一事件ID(雪花算法)
user_id BIGINT NOT NULL, -- 用户ID
event_type VARCHAR(32) NOT NULL, -- 事件类型
points_type TINYINT NOT NULL, -- 积分类型: 1=登录 2=充值 3=赠送
amount INT NOT NULL, -- 变动积分(正=入账, 负=出账)
balance_after INT NOT NULL, -- 该类型积分变动后余额(冗余快照)
biz_id VARCHAR(64) NOT NULL, -- 业务单号(幂等键)
biz_type VARCHAR(32) NOT NULL, -- 业务类型: recharge/consume/login/expire/refund
parent_event_id BIGINT DEFAULT NULL, -- 关联事件ID(冲正/退还时关联原事件)
metadata JSON DEFAULT NULL, -- 扩展信息
created_at DATETIME(3) NOT NULL, -- 毫秒级时间戳
INDEX idx_user_time (user_id, created_at),
UNIQUE KEY uk_biz_id (biz_id) -- 幂等约束
);
(2)积分快照表 points_snapshot¶
CREATE TABLE points_snapshot (
id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
login_points INT NOT NULL DEFAULT 0, -- 登录积分余额
login_expire_at DATETIME DEFAULT NULL, -- 登录积分过期时间
purchase_points INT NOT NULL DEFAULT 0, -- 充值积分余额
bonus_points INT NOT NULL DEFAULT 0, -- 赠送积分余额
last_event_id BIGINT NOT NULL, -- 快照对应的最后事件ID
snapshot_at DATETIME(3) NOT NULL, -- 快照时间
UNIQUE KEY uk_user (user_id),
INDEX idx_last_event (last_event_id)
);
(3)消耗流水表 consume_record¶
CREATE TABLE consume_record (
consume_id VARCHAR(64) PRIMARY KEY, -- 消耗单号
user_id BIGINT NOT NULL,
sub_account_id BIGINT DEFAULT NULL, -- 子账号ID
project_type VARCHAR(16) NOT NULL, -- 1K / 2K / video
total_cost INT NOT NULL, -- 总消耗积分
login_used INT NOT NULL DEFAULT 0,
purchase_used INT NOT NULL DEFAULT 0,
bonus_used INT NOT NULL DEFAULT 0,
order_id VARCHAR(64) DEFAULT NULL, -- 关联生成订单
status TINYINT NOT NULL DEFAULT 0, -- 0=处理中 1=成功 2=已回滚
idempotency_key VARCHAR(64) NOT NULL, -- 请求幂等键
created_at DATETIME(3) NOT NULL,
INDEX idx_user_time (user_id, created_at),
UNIQUE KEY uk_idempotency (idempotency_key)
);
(4)充值流水表 recharge_record¶
CREATE TABLE recharge_record (
recharge_id VARCHAR(64) PRIMARY KEY,
user_id BIGINT NOT NULL,
amount_yuan DECIMAL(10,2) NOT NULL, -- 实付金额(元)
purchase_points INT NOT NULL, -- 充值积分
bonus_points INT NOT NULL, -- 赠送积分
plan_id VARCHAR(32) NOT NULL, -- 套餐ID
payment_trade_no VARCHAR(64) DEFAULT NULL, -- 支付网关流水号
status TINYINT NOT NULL DEFAULT 0, -- 0=待支付 1=已支付 2=已退款
idempotency_key VARCHAR(64) NOT NULL,
created_at DATETIME(3) NOT NULL,
UNIQUE KEY uk_idempotency (idempotency_key)
);
(5)提成流水表 commission_record¶
CREATE TABLE commission_record (
commission_id VARCHAR(64) PRIMARY KEY,
user_id BIGINT NOT NULL, -- 消费者
sales_id BIGINT NOT NULL, -- 销售/代理
level VARCHAR(16) NOT NULL, -- sales / agentA / agentB
purchase_points_consumed INT NOT NULL, -- 消耗充值积分数
consume_id VARCHAR(64) NOT NULL, -- 关联消耗单号
rate DECIMAL(5,4) NOT NULL, -- 提成比例
commission_amount DECIMAL(10,2) NOT NULL, -- 提成金额
settle_month CHAR(7) NOT NULL, -- 结算月 '2026-07'
status TINYINT NOT NULL DEFAULT 0, -- 0=待结算 1=已结算 2=已支付
created_at DATETIME(3) NOT NULL,
UNIQUE KEY uk_consume_sales (consume_id, sales_id, level)
);
(6)会员状态表 membership¶
CREATE TABLE membership (
user_id BIGINT PRIMARY KEY,
plan_id VARCHAR(32) NOT NULL, -- 当前套餐
started_at DATETIME NOT NULL,
expires_at DATETIME NOT NULL,
auto_renew TINYINT DEFAULT 0,
sub_account_max INT NOT NULL,
concurrent_max INT NOT NULL,
storage_mb INT NOT NULL,
priority_gen TINYINT DEFAULT 0,
support_level VARCHAR(16) NOT NULL,
updated_at DATETIME(3) NOT NULL,
INDEX idx_expires (expires_at)
);
(7)操作审计表 audit_log¶
CREATE TABLE audit_log (
audit_id BIGINT PRIMARY KEY,
user_id BIGINT DEFAULT NULL,
action VARCHAR(64) NOT NULL, -- 操作类型
target_type VARCHAR(32) NOT NULL, -- 操作对象类型
target_id VARCHAR(64) NOT NULL, -- 操作对象ID
before_snapshot JSON DEFAULT NULL, -- 操作前数据快照
after_snapshot JSON DEFAULT NULL, -- 操作后数据快照
operator_id BIGINT DEFAULT NULL, -- 操作人
operator_ip VARCHAR(45) DEFAULT NULL,
trace_id VARCHAR(64) NOT NULL, -- 全链路追踪ID
created_at DATETIME(3) NOT NULL,
INDEX idx_user_time (user_id, created_at),
INDEX idx_target (target_type, target_id)
);
2.2 金额存储规范¶
┌──────────────────────────────────────────────────────┐
│ 铁律:所有金额字段使用 DECIMAL,禁止使用 FLOAT/DOUBLE │
├──────────────────────────────────────────────────────┤
│ 积分数量:INT(整数,无小数) │
│ 人民币金额:DECIMAL(10,2) │
│ 提成比例:DECIMAL(5,4),如 0.0300 = 3% │
│ 美元金额:DECIMAL(10,4),如成本核算 │
└──────────────────────────────────────────────────────┘
三、幂等性设计¶
3.1 三层幂等防护¶
请求到达
│
┌─────────▼─────────┐
│ Layer 1: 请求级 │ 客户端生成 idempotency_key
│ Redis 去重缓存 │ TTL=24h,SET NX 原子操作
└─────────┬─────────┘
│ 通过
┌─────────▼─────────┐
│ Layer 2: 业务级 │ 数据库唯一索引 uk_biz_id
│ DB唯一约束 │ 插入重复抛 DuplicateKeyException
└─────────┬─────────┘
│ 通过
┌─────────▼─────────┐
│ Layer 3: 事件级 │ 查询已有事件,直接返回成功
│ 查重后返回 │ "已处理,无需重复执行"
└────────────────────┘
3.2 幂等键生成规范¶
| 业务场景 | 幂等键格式 | 示例 |
|---|---|---|
| 积分消耗 | CONSUME:{order_id} |
CONSUME:ORD20260721001 |
| 充值到账 | RECHARGE:{trade_no} |
RECHARGE:4200001234567890 |
| 登录积分 | LOGIN:{user_id}:{date} |
LOGIN:10086:2026-07-21 |
| 积分过期 | EXPIRE:{user_id}:{date} |
EXPIRE:10086:2026-07-21 |
| 提成计算 | COMM:{consume_id}:{sales_id} |
COMM:C20260721001:2001 |
3.3 实现示例(伪代码)¶
async def consume_points(user_id: int, project_type: str, order_id: str):
idempotency_key = f"CONSUME:{order_id}"
# Layer 1: Redis 去重
if not await redis.set(idempotency_key, "processing", nx=True, ex=86400):
existing = await redis.get(idempotency_key)
if existing == "done":
return {"status": "already_done", "consume_id": await get_consume_id(order_id)}
else:
raise HttpException(409, "请求处理中,请勿重复提交")
try:
# Layer 2 & 3: 数据库层
async with db.transaction(isolation="SERIALIZABLE"):
# 查重
existing = await db.query_one(
"SELECT * FROM consume_record WHERE idempotency_key = %s",
idempotency_key
)
if existing:
await redis.set(idempotency_key, "done", ex=86400)
return {"status": "already_done", "consume_id": existing.consume_id}
# 执行消耗逻辑 ...
result = await _execute_consume(user_id, project_type, order_id)
await redis.set(idempotency_key, "done", ex=86400)
return result
except Exception as e:
await redis.delete(idempotency_key) # 失败则释放,允许重试
raise
四、事务与并发控制¶
4.1 并发场景分析¶
| 场景 | 风险 | 方案 |
|---|---|---|
| 同一用户同时消耗 | 超扣积分 | 行级锁 + 悲观锁 |
| 充值到账与消耗并发 | 余额计算错乱 | 事件队列串行化 |
| 登录积分同时领取 | 重复领取 | 日期唯一索引 + Redis锁 |
| 多子账号同时消耗 | 共享积分池超扣 | 主账号级别分布式锁 |
| 提成计算与消耗并发 | 提成金额不准 | 消耗确认后再算提成(最终一致性) |
4.2 积分消耗的原子操作¶
async def _execute_consume(user_id: int, project_type: str, order_id: str):
"""
核心消耗逻辑 —— 必须在一个数据库事务内完成
"""
cost = POINTS_COST[project_type] # 1K=12, 2K=18, video=84/10s
async with db.transaction(isolation="REPEATABLE_READ") as tx:
# 1. 悲观锁锁定用户积分快照行
snapshot = await db.query_one(
"""SELECT * FROM points_snapshot
WHERE user_id = %s FOR UPDATE""",
user_id
)
if not snapshot:
raise Exception("用户积分账户不存在")
# 2. 计算各类型积分可扣减量
now = datetime.now()
# 登录积分检查过期
login_available = snapshot.login_points
if snapshot.login_expire_at and snapshot.login_expire_at < now:
login_available = 0 # 已过期,不可用
remaining = cost
login_deduct = min(remaining, login_available)
remaining -= login_deduct
purchase_deduct = min(remaining, snapshot.purchase_points)
remaining -= purchase_deduct
bonus_deduct = min(remaining, snapshot.bonus_points)
remaining -= bonus_deduct
if remaining > 0:
raise InsufficientPointsException(f"积分不足,还差 {remaining} 积分")
# 3. 写入事件(不可变表,只追加)
events_to_insert = []
if login_deduct > 0:
events_to_insert.append((user_id, 'consume', 1, -login_deduct,
snapshot.login_points - login_deduct,
f"CONSUME:{order_id}", 'consume'))
snapshot.login_points -= login_deduct
if purchase_deduct > 0:
events_to_insert.append((user_id, 'consume', 2, -purchase_deduct,
snapshot.purchase_points - purchase_deduct,
f"CONSUME:{order_id}", 'consume'))
snapshot.purchase_points -= purchase_deduct
if bonus_deduct > 0:
events_to_insert.append((user_id, 'consume', 3, -bonus_deduct,
snapshot.bonus_points - bonus_deduct,
f"CONSUME:{order_id}", 'consume'))
snapshot.bonus_points -= bonus_deduct
# 批量写入事件
event_ids = []
for evt in events_to_insert:
event_id = snowflake.next_id()
await db.execute(
"""INSERT INTO points_event
(event_id, user_id, event_type, points_type, amount,
balance_after, biz_id, biz_type, created_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
event_id, *evt, now
)
event_ids.append(event_id)
# 4. 更新快照
await db.execute(
"""UPDATE points_snapshot
SET login_points = %s, purchase_points = %s, bonus_points = %s,
last_event_id = %s, snapshot_at = %s
WHERE user_id = %s""",
snapshot.login_points, snapshot.purchase_points, snapshot.bonus_points,
event_ids[-1], now, user_id
)
# 5. 写入消耗流水
consume_id = f"C{now.strftime('%Y%m%d%H%M%S')}{user_id}{random_int(4)}"
await db.execute(
"""INSERT INTO consume_record
(consume_id, user_id, project_type, total_cost,
login_used, purchase_used, bonus_used, order_id,
status, idempotency_key, created_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,1,%s,%s)""",
consume_id, user_id, project_type, cost,
login_deduct, purchase_deduct, bonus_deduct,
order_id, f"CONSUME:{order_id}", now
)
return {
"consume_id": consume_id,
"total_cost": cost,
"login_used": login_deduct,
"purchase_used": purchase_deduct,
"bonus_used": bonus_deduct,
"remaining": {
"login_points": snapshot.login_points,
"purchase_points": snapshot.purchase_points,
"bonus_points": snapshot.bonus_points
}
}
4.3 共享积分池的分布式锁(子账号场景)¶
async def consume_by_subaccount(main_user_id: int, sub_account_id: int,
project_type: str, order_id: str):
"""
子账号消耗 —— 需要主账号级别的分布式锁
"""
lock_key = f"lock:points:{main_user_id}"
lock = redis.lock(lock_key, timeout=10) # 10秒超时
acquired = await lock.acquire(blocking=True, blocking_timeout=5)
if not acquired:
raise HttpException(429, "操作频繁,请稍后重试")
try:
result = await _execute_consume(main_user_id, project_type, order_id)
# 在消耗流水中标记子账号
await db.execute(
"UPDATE consume_record SET sub_account_id = %s WHERE consume_id = %s",
sub_account_id, result["consume_id"]
)
return result
finally:
await lock.release()
五、审计与可回溯¶
5.1 不可变性(Immutability)¶
原则:已确认的流水记录永不UPDATE,只INSERT
┌─────────────────────────────────────────────────────┐
│ 正常操作:INSERT 新事件 │
│ 错误修正:INSERT 冲正事件(红字冲销),而非DELETE原事件 │
│ 退款:INSERT 退还事件,关联原事件 │
│ 过期:INSERT 过期事件,系统自动触发 │
└─────────────────────────────────────────────────────┘
5.2 冲正(红字冲销)示例¶
async def reverse_consume(original_consume_id: str, reason: str, operator_id: int):
"""
冲正一笔消耗 —— 不删除原记录,而是插入反向事件
"""
async with db.transaction():
# 1. 查原消耗记录
original = await db.query_one(
"SELECT * FROM consume_record WHERE consume_id = %s FOR UPDATE",
original_consume_id
)
if not original:
raise Exception("原消耗记录不存在")
if original.status == 2:
raise Exception("该记录已被冲正")
# 2. 插入反向事件(积分返还)
events = []
if original.login_used > 0:
events.append((original.user_id, 'reverse', 1, original.login_used,
f"REVERSE:{original_consume_id}", 'reverse'))
if original.purchase_used > 0:
events.append((original.user_id, 'reverse', 2, original.purchase_used,
f"REVERSE:{original_consume_id}", 'reverse'))
if original.bonus_used > 0:
events.append((original.user_id, 'reverse', 3, original.bonus_used,
f"REVERSE:{original_consume_id}", 'reverse'))
for evt in events:
await db.execute(
"""INSERT INTO points_event
(event_id, user_id, event_type, points_type, amount,
balance_after, biz_id, biz_type, parent_event_id, metadata, created_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
snowflake.next_id(), *evt,
json.dumps({"reason": reason, "operator_id": operator_id}),
datetime.now()
)
# 3. 更新原消耗记录状态
await db.execute(
"UPDATE consume_record SET status = 2 WHERE consume_id = %s",
original_consume_id
)
# 4. 重新计算快照
await rebuild_snapshot(original.user_id)
5.3 审计日志自动记录¶
# 使用数据库触发器或ORM钩子自动记录关键操作
# 示例:MySQL触发器
"""
CREATE TRIGGER trg_points_event_audit
AFTER INSERT ON points_event
FOR EACH ROW
BEGIN
INSERT INTO audit_log (audit_id, user_id, action, target_type, target_id,
after_snapshot, trace_id, created_at)
VALUES (NEXT_ID(), NEW.user_id, NEW.event_type, 'points_event', NEW.event_id,
JSON_OBJECT('points_type', NEW.points_type, 'amount', NEW.amount,
'balance_after', NEW.balance_after),
@trace_id, NOW());
END;
"""
5.4 全链路追踪¶
每个请求携带 trace_id,贯穿:
网关 → 业务服务 → 数据库 → 日志 → 审计表
实现方式:
- HTTP Header: X-Trace-Id
- 服务间透传(gRPC metadata / MQ header)
- MySQL 会话变量: SET @trace_id = 'xxx'
- 日志格式: [trace_id] [user_id] [action] message
六、日终对账体系¶
6.1 对账公式¶
快照余额 vs 事件累计 = 一致性校验
对任意用户 user_id:
snapshot.login_points ?= SUM(points_event.amount WHERE points_type=1 AND user_id=user_id)
snapshot.purchase_points ?= SUM(points_event.amount WHERE points_type=2 AND user_id=user_id)
snapshot.bonus_points ?= SUM(points_event.amount WHERE points_type=3 AND user_id=user_id)
整体平台:
总充值积分(T0时刻) + Σ入账事件 - Σ出账事件 = 总余额(T1时刻)
Σ(consume_record.total_cost) = Σ(points_event.amount WHERE amount<0 AND event_type!='reverse')
6.2 日终对账任务¶
async def daily_reconciliation(target_date: str):
"""
每日凌晨执行,对账流程
"""
results = {
"date": target_date,
"total_users": 0,
"matched": 0,
"mismatched": [],
"errors": []
}
# 1. 获取所有用户快照
snapshots = await db.query_all(
"""SELECT * FROM points_snapshot
WHERE snapshot_at >= %s AND snapshot_at < %s""",
f"{target_date} 00:00:00", f"{target_date} 23:59:59"
)
results["total_users"] = len(snapshots)
# 2. 逐用户校验
for snap in snapshots:
# 从事件表聚合余额
event_summary = await db.query_one(
"""SELECT
SUM(CASE WHEN points_type=1 THEN amount ELSE 0 END) as login_sum,
SUM(CASE WHEN points_type=2 THEN amount ELSE 0 END) as purchase_sum,
SUM(CASE WHEN points_type=3 THEN amount ELSE 0 END) as bonus_sum
FROM points_event
WHERE user_id = %s""",
snap.user_id
)
diffs = []
if snap.login_points != (event_summary.login_sum or 0):
diffs.append(f"登录积分: 快照={snap.login_points} 事件累计={event_summary.login_sum}")
if snap.purchase_points != (event_summary.purchase_sum or 0):
diffs.append(f"充值积分: 快照={snap.purchase_points} 事件累计={event_summary.purchase_sum}")
if snap.bonus_points != (event_summary.bonus_sum or 0):
diffs.append(f"赠送积分: 快照={snap.bonus_points} 事件累计={event_summary.bonus_sum}")
if diffs:
results["mismatched"].append({
"user_id": snap.user_id,
"diffs": diffs,
"last_event_id": snap.last_event_id
})
# 自动修复
await rebuild_snapshot(snap.user_id)
else:
results["matched"] += 1
# 3. 消耗总额校验
consume_check = await db.query_one(
"""SELECT
(SELECT SUM(total_cost) FROM consume_record
WHERE DATE(created_at)=%s AND status=1) as consume_total,
(SELECT SUM(ABS(amount)) FROM points_event
WHERE DATE(created_at)=%s AND event_type='consume') as event_total
""", target_date, target_date
)
if consume_check.consume_total != consume_check.event_total:
results["errors"].append(
f"消耗总额不一致: 流水={consume_check.consume_total} 事件={consume_check.event_total}"
)
# 4. 生成对账报告
await save_reconciliation_report(target_date, results)
# 5. 如有差异,触发告警
if results["mismatched"] or results["errors"]:
await alert_service.send_alert(
level="CRITICAL",
title=f"日终对账异常 - {target_date}",
detail=json.dumps(results, ensure_ascii=False)
)
return results
6.3 快照重建(自动修复)¶
async def rebuild_snapshot(user_id: int):
"""
基于事件表全量重建用户快照 —— 对账异常的修复手段
"""
async with db.transaction():
# 锁定用户行
await db.execute(
"SELECT * FROM points_snapshot WHERE user_id = %s FOR UPDATE",
user_id
)
# 聚合所有事件
summary = await db.query_one(
"""SELECT
SUM(CASE WHEN points_type=1 THEN amount ELSE 0 END) as login_sum,
SUM(CASE WHEN points_type=2 THEN amount ELSE 0 END) as purchase_sum,
SUM(CASE WHEN points_type=3 THEN amount ELSE 0 END) as bonus_sum,
MAX(event_id) as last_event_id
FROM points_event
WHERE user_id = %s""",
user_id
)
# 获取登录积分过期时间
login_expire = await db.query_one(
"""SELECT DATE_ADD(created_at, INTERVAL 24 HOUR) as expire_at
FROM points_event
WHERE user_id = %s AND points_type=1 AND event_type='login_grant'
ORDER BY event_id DESC LIMIT 1""",
user_id
)
# UPSERT 快照
await db.execute(
"""INSERT INTO points_snapshot
(user_id, login_points, login_expire_at, purchase_points, bonus_points,
last_event_id, snapshot_at)
VALUES (%s,%s,%s,%s,%s,%s,NOW())
ON DUPLICATE KEY UPDATE
login_points=VALUES(login_points),
login_expire_at=VALUES(login_expire_at),
purchase_points=VALUES(purchase_points),
bonus_points=VALUES(bonus_points),
last_event_id=VALUES(last_event_id),
snapshot_at=NOW()""",
user_id,
summary.login_sum or 0,
login_expire.expire_at if login_expire else None,
summary.purchase_sum or 0,
summary.bonus_sum or 0,
summary.last_event_id or 0
)
七、持久性与容灾¶
7.1 数据库层¶
┌─────────────────────────────────────────────────────┐
│ 主库(写入) │
│ │ │
│ ├── 半同步复制 ──► 从库1(读取 + 灾备) │
│ │ │
│ └── 异步复制 ────► 从库2(报表查询 + 对账) │
│ │
│ Binlog 开启 ROW 格式 + GTID │
│ 每日全量备份 + 每小时增量备份 │
│ 备份保留:近7天每日 + 近12周每周 + 近12月每月 │
└─────────────────────────────────────────────────────┘
7.2 WAL(Write-Ahead Logging)¶
MySQL InnoDB 默认 WAL 机制:
- 写操作先写 redo log(顺序写,极快)
- 后异步刷盘到数据文件
- 崩溃恢复时从 redo log 重放
建议配置:
innodb_flush_log_at_trx_commit = 1 # 每次事务都刷盘(最强持久性)
sync_binlog = 1 # 每次事务同步binlog
innodb_doublewrite = ON # 防止页断裂
7.3 多级备份策略¶
| 层级 | 内容 | 频率 | 保留 |
|---|---|---|---|
| Binlog | 增量变更 | 实时 | 30天 |
| 全量备份 | mysqldump | 每日凌晨3点 | 30天 |
| 冷备份 | 异地OSS | 每周 | 12周 |
| 归档 | 年度归档 | 每年 | 永久 |
7.4 事件溯源天然优势¶
事件表本身就是完整的数据源:
- 即使快照表全部损坏,从事件表可完全重建
- 事件表配合 binlog 可恢复到任意时间点
- 天然支持审计和时间旅行查询
八、健壮性保障¶
8.1 重试机制¶
@retry(
max_attempts=3,
backoff=exponential_backoff(base=1, factor=2), # 1s → 2s → 4s
retry_on=[DeadlockError, ConnectionError, TimeoutError],
on_giveup=lambda e: alert_service.send(e)
)
async def safe_consume(user_id: int, project_type: str, order_id: str):
return await consume_points(user_id, project_type, order_id)
8.2 死信队列(DLQ)¶
积分操作失败时:
正常 → 重试3次 → 仍失败 → 写入死信队列 → 人工介入
死信队列记录:
- 原始请求参数
- 失败原因
- 失败时间
- 重试次数
人工处理工具:从死信队列重放或取消
8.3 熔断降级¶
class PointsCircuitBreaker:
"""
当数据库异常率超过阈值时自动熔断
"""
def __init__(self):
self.failure_threshold = 0.5 # 50% 失败率
self.half_open_timeout = 60 # 60秒后半开
self.window_size = 100 # 滑动窗口大小
self.state = "CLOSED" # CLOSED / OPEN / HALF_OPEN
async def call(self, func, *args):
if self.state == "OPEN":
if time_since_last_failure() > self.half_open_timeout:
self.state = "HALF_OPEN"
else:
raise CircuitBreakerOpenException("积分服务熔断中,请稍后再试")
try:
result = await func(*args)
self._record_success()
return result
except Exception:
self._record_failure()
if self.state == "HALF_OPEN":
self.state = "OPEN"
elif self.failure_rate() > self.failure_threshold:
self.state = "OPEN"
alert_service.send("积分服务已熔断")
raise
8.4 限流¶
# 单用户消耗限流(防止恶意刷量)
user_rate_limiter = TokenBucket(
key="rate:consume:{user_id}",
rate=10, # 每秒最多10次
burst=20 # 突发允许20次
)
# 全局限流
global_rate_limiter = TokenBucket(
key="rate:consume:global",
rate=1000, # 全局每秒1000次
burst=2000
)
8.5 登录积分过期处理¶
async def expire_login_points():
"""
定时任务:每分钟扫描过期登录积分
不依赖定时任务做精确过期判断,消耗时实时检查过期时间
"""
now = datetime.now()
# 批量查询已过期的登录积分
expired = await db.query_all(
"""SELECT user_id, login_points, login_expire_at
FROM points_snapshot
WHERE login_points > 0 AND login_expire_at <= %s
LIMIT 1000""",
now
)
for row in expired:
async with db.transaction():
# 锁定行
snap = await db.query_one(
"SELECT * FROM points_snapshot WHERE user_id = %s FOR UPDATE",
row.user_id
)
# 二次确认(防止并发修改)
if snap.login_points > 0 and snap.login_expire_at and snap.login_expire_at <= now:
expire_amount = snap.login_points
# 写入过期事件
await db.execute(
"""INSERT INTO points_event
(event_id, user_id, event_type, points_type, amount,
balance_after, biz_id, biz_type, created_at balance_after, biz_id, biz_type, created_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
snowflake.next_id(), row.user_id, 'expire', 1, -expire_amount,
0, f"EXPIRE:{row.user_id}:{now.strftime('%Y%m%d')}", 'expire', now
)
# 更新快照
await db.execute(
"""UPDATE points_snapshot
SET login_points = 0, login_expire_at = NULL, snapshot_at = NOW()
WHERE user_id = %s""",
row.user_id
)
### 8.6 会员到期自动降级
```python
async def downgrade_expired_memberships():
"""
每日凌晨执行,将所有过期会员降级为体验档
"""
now = datetime.now()
expired = await db.query_all(
"""SELECT user_id, plan_id FROM membership
WHERE expires_at <= %s AND plan_id != 'trial'
LIMIT 500""",
now
)
for row in expired:
async with db.transaction():
await db.execute(
"""UPDATE membership
SET plan_id = 'trial',
sub_account_max = 0, concurrent_max = 1,
storage_mb = 1024, priority_gen = 0,
support_level = 'normal', updated_at = NOW()
WHERE user_id = %s AND expires_at <= %s""",
row.user_id, now
)
# 记录审计
await db.execute(
"""INSERT INTO audit_log
(audit_id, user_id, action, target_type, target_id,
before_snapshot, after_snapshot, trace_id, created_at)
VALUES (%s,%s,'membership_expire','membership',%s,
%s,%s,'system',NOW())""",
snowflake.next_id(), row.user_id, row.user_id,
json.dumps({"plan": row.plan_id}), json.dumps({"plan": "trial"})
)
九、关键业务流程¶
9.1 积分消耗完整流程(时序图)¶
sequenceDiagram
participant C as 客户端
participant G as API网关
participant S as 积分服务
participant R as Redis
participant DB as MySQL
participant MQ as 消息队列
C->>G: POST /api/points/consume
G->>G: 校验鉴权 + 限流
G->>S: consume_points(user_id, type, order_id)
S->>R: SET NX idempotency_key (幂等L1)
alt Key已存在
R-->>S: 已存在
S-->>C: 409 重复请求
end
S->>DB: BEGIN TRANSACTION
S->>DB: SELECT ... FOR UPDATE (锁定快照行)
S->>S: 计算消耗优先级: 登录→充值→赠送
alt 积分不足
S->>DB: ROLLBACK
S->>R: DEL idempotency_key
S-->>C: 402 积分不足
end
S->>DB: INSERT points_event × N (不可变事件)
S->>DB: UPDATE points_snapshot (刷新快照)
S->>DB: INSERT consume_record (消耗流水)
S->>DB: COMMIT
S->>MQ: 发送 consume_event (提成异步计算)
S->>R: SET idempotency_key = "done"
S-->>C: 200 消耗成功 + 剩余余额
MQ->>S: 异步消费 → 计算提成
S->>DB: INSERT commission_record
9.2 充值流程¶
async def recharge(user_id: int, plan_id: str, payment_trade_no: str):
"""
充值到账流程 —— 支付回调触发
"""
idempotency_key = f"RECHARGE:{payment_trade_no}"
# 幂等检查
existing = await db.query_one(
"SELECT * FROM recharge_record WHERE idempotency_key = %s",
idempotency_key
)
if existing:
return {"status": "already_processed", "recharge_id": existing.recharge_id}
# 查套餐配置
plan = PLANS[plan_id]
purchase_points = int(plan.price * 10) # 1元=10积分
bonus_points = int(purchase_points * plan.bonus_rate)
async with db.transaction():
# 1. 写充值流水
recharge_id = f"R{datetime.now().strftime('%Y%m%d%H%M%S')}{user_id}"
await db.execute(
"""INSERT INTO recharge_record
(recharge_id, user_id, amount_yuan, purchase_points, bonus_points,
plan_id, payment_trade_no, status, idempotency_key, created_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,1,%s,NOW())""",
recharge_id, user_id, plan.price, purchase_points, bonus_points,
plan_id, payment_trade_no, idempotency_key
)
# 2. 写积分事件:充值积分入账
snap = await db.query_one(
"SELECT * FROM points_snapshot WHERE user_id = %s FOR UPDATE", user_id
)
new_purchase = (snap.purchase_points if snap else 0) + purchase_points
await db.execute(
"""INSERT INTO points_event
(event_id, user_id, event_type, points_type, amount,
balance_after, biz_id, biz_type, created_at)
VALUES (%s,%s,'recharge',2,%s,%s,%s,'recharge',NOW())""",
snowflake.next_id(), user_id, purchase_points,
new_purchase, recharge_id
)
# 3. 写积分事件:赠送积分入账
if bonus_points > 0:
new_bonus = (snap.bonus_points if snap else 0) + bonus_points
await db.execute(
"""INSERT INTO points_event
(event_id, user_id, event_type, points_type, amount,
balance_after, biz_id, biz_type, created_at)
VALUES (%s,%s,'bonus_grant',3,%s,%s,%s,'recharge',NOW())""",
snowflake.next_id(), user_id, bonus_points,
new_bonus, recharge_id
)
# 4. 更新快照
period_days = PLAN_PERIOD_DAYS[plan.period] # 月=30 季=90 年=365
await db.execute(
"""INSERT INTO points_snapshot
(user_id, login_points, purchase_points, bonus_points,
last_event_id, snapshot_at)
VALUES (%s,0,%s,%s,%s,NOW())
ON DUPLICATE KEY UPDATE
purchase_points = purchase_points + VALUES(purchase_points),
bonus_points = bonus_points + VALUES(bonus_points),
last_event_id = VALUES(last_event_id),
snapshot_at = NOW()""",
user_id, purchase_points, bonus_points, snowflake.last_id()
)
# 5. 更新/创建会员
await db.execute(
"""INSERT INTO membership
(user_id, plan_id, started_at, expires_at,
sub_account_max, concurrent_max, storage_mb,
priority_gen, support_level, updated_at)
VALUES (%s,%s,NOW(),DATE_ADD(NOW(), INTERVAL %s DAY),
%s,%s,%s,%s,%s,NOW())
ON DUPLICATE KEY UPDATE
plan_id = VALUES(plan_id),
started_at = VALUES(started_at),
expires_at = VALUES(expires_at),
sub_account_max = VALUES(sub_account_max),
concurrent_max = VALUES(concurrent_max),
storage_mb = VALUES(storage_mb),
priority_gen = VALUES(priority_gen),
support_level = VALUES(support_level),
updated_at = NOW()""",
user_id, plan_id, period_days,
plan.sub_accounts, plan.concurrent, plan.storage,
1 if plan.priority else 0, plan.support_level
)
return {
"recharge_id": recharge_id,
"purchase_points": purchase_points,
"bonus_points": bonus_points,
"total_points": purchase_points + bonus_points
}
9.3 登录积分领取¶
async def claim_login_points(user_id: int, claim_date: str):
"""
每日登录积分领取
claim_date 格式: YYYY-MM-DD
"""
idempotency_key = f"LOGIN:{user_id}:{claim_date}"
# Layer 1: Redis
if not await redis.set(idempotency_key, "processing", nx=True, ex=86400):
raise HttpException(409, "今日已领取")
try:
async with db.transaction():
# 查重
existing = await db.query_one(
"SELECT event_id FROM points_event WHERE biz_id = %s",
idempotency_key
)
if existing:
raise HttpException(409, "今日已领取")
# 查会员状态
member = await db.query_one(
"SELECT * FROM membership WHERE user_id = %s AND expires_at > NOW()",
user_id
)
if not member:
raise HttpException(403, "会员已过期,无法领取登录积分")
plan = PLANS[member.plan_id]
if plan.login_points_per_day <= 0:
raise HttpException(403, "当前套餐无登录积分")
# 锁定快照
snap = await db.query_one(
"SELECT * FROM points_snapshot WHERE user_id = %s FOR UPDATE",
user_id
)
# 写入入账事件
new_login = (snap.login_points if snap else 0) + plan.login_points_per_day
expire_at = datetime.now() + timedelta(hours=24)
await db.execute(
"""INSERT INTO points_event
(event_id, user_id, event_type, points_type, amount,
balance_after, biz_id, biz_type, created_at)
VALUES (%s,%s,'login_grant',1,%s,%s,%s,'login',NOW())""",
snowflake.next_id(), user_id, plan.login_points_per_day,
new_login, idempotency_key
)
# 更新快照(含过期时间)
await db.execute(
"""UPDATE points_snapshot
SET login_points = %s, login_expire_at = %s,
snapshot_at = NOW()
WHERE user_id = %s""",
new_login, expire_at, user_id
)
await redis.set(idempotency_key, "done", ex=86400)
return {
"granted": plan.login_points_per_day,
"login_points": new_login,
"expire_at": expire_at.isoformat()
}
except HttpException:
await redis.delete(idempotency_key)
raise
9.4 提成异步计算¶
async def calculate_commission(consume_id: str):
"""
消费事件驱动:消耗发生后异步计算提成
通过消息队列触发
"""
# 查消耗记录
consume = await db.query_one(
"SELECT * FROM consume_record WHERE consume_id = %s AND status = 1",
consume_id
)
if not consume or consume.purchase_used <= 0:
return # 没有消耗充值积分,不提成
# 查该用户的销售/代理关系
relations = await db.query_all(
"""SELECT sales_id, level FROM user_sales_relation
WHERE user_id = %s""",
consume.user_id
)
now = datetime.now()
settle_month = now.strftime('%Y-%m')
for rel in relations:
rate = COMMISSION_RATES[rel.level] # sales=3%, agentA=12%, agentB=8%
amount = Decimal(consume.purchase_used) / Decimal(10) * rate
commission_id = f"COMM{now.strftime('%Y%m%d%H%M%S')}{consume.user_id}{rel.sales_id}"
# 幂等:consume_id + sales_id + level 唯一
await db.execute(
"""INSERT IGNORE INTO commission_record
(commission_id, user_id, sales_id, level, purchase_points_consumed,
consume_id, rate, commission_amount, settle_month, status, created_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,0,NOW())""",
commission_id, consume.user_id, rel.sales_id, rel.level,
consume.purchase_used, consume_id, rate, amount, settle_month
)
9.5 提成月结算¶
async def monthly_commission_settlement(month: str):
"""
每月1日执行上个月提成结算
"""
async with db.transaction():
# 锁定当月所有待结算记录
records = await db.query_all(
"""SELECT * FROM commission_record
WHERE settle_month = %s AND status = 0
FOR UPDATE""",
month
)
# 按销售汇总
summary = {}
for r in records:
key = (r.sales_id, r.level)
if key not in summary:
summary[key] = {"total": Decimal(0), "details": []}
summary[key]["total"] += r.commission_amount
summary[key]["details"].append(r.commission_id)
# 更新为已结算
await db.execute(
"""UPDATE commission_record
SET status = 1 WHERE settle_month = %s AND status = 0""",
month
)
# 生成结算报表
for (sales_id, level), data in summary.items():
await db.execute(
"""INSERT INTO commission_settlement
(settlement_id, sales_id, level, month, total_amount, status, created_at)
VALUES (%s,%s,%s,%s,%s,'pending',NOW())""",
snowflake.next_id(), sales_id, level, month, data["total"]
)
return summary
十、监控与告警¶
10.1 核心监控指标¶
| 指标 | 类型 | 告警阈值 | 说明 |
|---|---|---|---|
| 消耗成功率 | Gauge | < 99.5% | 过去5分钟消耗请求成功率 |
| 消耗P99延迟 | Histogram | > 500ms | 消耗接口P99响应时间 |
| 日终对账差异数 | Counter | > 0 | 对账发现差异立即告警 |
| 死信队列堆积 | Gauge | > 50 | 死信队列待处理数量 |
| 数据库连接池使用率 | Gauge | > 80% | 连接池接近耗尽 |
| 会员过期处理延迟 | Gauge | > 300s | 过期处理任务延迟 |
| 积分事件写入速率 | Meter | 异常波动 > 50% | 突增或骤降 |
| 日消耗积分总额 | Meter | 日环比波动 > 100% | 异常消耗监测 |
10.2 业务监控大盘¶
┌─────────────────────────────────────────────────────────┐
│ 📊 积分系统业务大盘 │
├───────────────┬───────────────┬─────────────────────────┤
│ 今日消耗总积分 │ 今日充值总额 │ 今日登录积分领取人次 │
│ 1,234,567 │ ¥156,789.00 │ 3,456 │
├───────────────┼───────────────┼─────────────────────────┤
│ 消耗成功率 │ P99延迟 │ 待处理死信 │
│ 99.87% ✅ │ 187ms ✅ │ 0 ✅ │
├───────────────┴───────────────┴─────────────────────────┤
│ 日终对账状态: ✅ 2026-07-20 对账通过 (100,000/100,000) │
│ 最近异常: 2026-07-18 用户10086 快照不一致 (已自动修复) │
└─────────────────────────────────────────────────────────┘
10.3 关键告警规则¶
ALERT_RULES = [
{
"name": "对账差异告警",
"condition": "reconciliation_mismatch_count > 0",
"level": "P0 - 紧急",
"channel": ["电话", "企业微信", "飞书"],
"message_template": "【P0】日终对账发现 {count} 个用户余额不一致,请立即排查!"
},
{
"name": "消耗成功率下降",
"condition": "consume_success_rate < 0.99 FOR 5m",
"level": "P1 - 严重",
"channel": ["企业微信", "飞书"],
"message_template": "【P1】积分消耗成功率降至 {rate}%,可能影响用户使用"
},
{
"name": "死信队列堆积",
"condition": "dlq_pending_count > 50 FOR 5m",
"level": "P1 - 严重",
"channel": ["企业微信", "飞书"],
"message_template": "【P1】死信队列堆积 {count} 条,需人工介入"
},
{
"name": "异常大额消耗",
"condition": "single_user_hourly_consume > 10000",
"level": "P2 - 关注",
"channel": ["飞书"],
"message_template": "【P2】用户 {user_id} 1小时内消耗 {amount} 积分,请关注"
},
{
"name": "数据库主从延迟",
"condition": "replication_lag > 10s FOR 1m",
"level": "P2 - 关注",
"channel": ["飞书"],
"message_template": "【P2】主从复制延迟 {lag}s,可能影响对账"
}
]
10.4 日志规范¶
日志级别规范:
ERROR - 数据库异常、事务失败、对账不一致(需立即处理)
WARN - 重试成功、限流触发、熔断半开、余额不足
INFO - 每次消耗/充值/登录积分的完整请求和响应
DEBUG - 完整SQL、Redis操作(仅开发环境)
关键 INFO 日志示例:
[INFO] [trace_id=abc123] [user_id=10086] [action=consume]
project_type=1K total_cost=12 login_used=12 purchase_used=0 bonus_used=0
remaining={login:0, purchase:5880, bonus:1280} elapsed=23ms
数据保留:
- 应用日志:30天(热)+ 90天(冷归档)
- 审计日志:永久保留
- 积分事件:永久保留(不可删除)
十一、技术选型建议¶
11.1 推荐技术栈¶
┌─────────────────────────────────────────────────────┐
│ 层级 │ 推荐技术 │ 备选 │
├─────────────────────────────────────────────────────┤
│ 语言 │ Python 3.12+ │ Go / Java │
│ Web框架 │ FastAPI │ Django │
│ ORM │ SQLAlchemy 2.0 │ — │
│ 数据库 │ MySQL 8.0+ (InnoDB) │ PostgreSQL │
│ 缓存 │ Redis 7.0+ │ — │
│ 消息队列 │ RabbitMQ / Kafka │ Redis Stream │
│ 定时任务 │ APScheduler │ Celery Beat │
│ 监控 │ Prometheus + Grafana│ — │
│ 日志 │ ELK (Elasticsearch) │ Loki │
│ 分布式ID │ Snowflake (自建) │ 美团Leaf │
│ Trace │ OpenTelemetry │ Jaeger │
│ 配置中心 │ Nacos / Consul │ — │
└─────────────────────────────────────────────────────┘
11.2 为什么选择事件溯源而非传统CRUD¶
| 维度 | 事件溯源 | 传统余额UPDATE |
|---|---|---|
| 可追溯性 | ✅ 完整事件链 | ❌ 只知道当前余额 |
| 审计 | ✅ 天然支持 | ❌ 需要额外审计表 |
| 对账 | ✅ 事件重放与快照互相校验 | ❌ 余额可能被覆盖 |
| 错误修正 | ✅ 插入冲正事件 | ⚠️ 直接修改余额(危险) |
| 查询性能 | ⚠️ 需快照加速 | ✅ 直接读余额 |
| 存储开销 | ⚠️ 事件表数据量大 | ✅ 只有当前余额 |
| 并发控制 | ✅ 事件追加无冲突 | ⚠️ 需悲观锁 |
结论:积分系统对准确性和可追溯性要求极高,事件溯源的开销完全值得。配合快照表解决了查询性能问题。
11.3 部署架构建议¶
┌──────────────┐
│ CDN / LB │
└──────┬───────┘
│
┌────────────┼────────────┐
│ │ │
┌──────▼──────┐ ┌──▼───┐ ┌─────▼──────┐
│ API Gateway │ │ API │ │ API │
│ (Kong/Nginx)│ │ GW-2 │ │ GW-3 │
└──────┬──────┘ └──┬───┘ └─────┬──────┘
│ │ │
└────────────┼────────────┘
│
┌────────────▼────────────┐
│ Points Service │
│ (多副本,无状态) │
└────────────┬────────────┘
│
┌─────────────────┼─────────────────┐
│ │ │
┌──────▼──────┐ ┌──────▼──────┐ ┌───────▼──────┐
│ MySQL │ │ Redis │ │ RabbitMQ │
│ (主从) │ │ (Cluster) │ │ (Cluster) │
└─────────────┘ └─────────────┘ └──────────────┘
十二、检查清单(上线前必检)¶
□ 所有金额字段使用 DECIMAL,无 FLOAT
□ 所有流水表有唯一幂等键约束
□ 消耗/充值核心逻辑在数据库事务内完成
□ SELECT ... FOR UPDATE 正确使用,无死锁风险
□ 事件表只追加不修改
□ 冲正操作为 INSERT 反向事件,非 DELETE
□ 日终对账定时任务已配置并测试通过
□ 快照重建脚本可用于紧急修复
□ 监控告警规则已配置(对账差异=CRITICAL)
□ 数据库备份策略已实施(全量+增量+异地)
□ Binlog ROW格式 + GTID 已开启
□ innodb_flush_log_at_trx_commit = 1
□ 登录积分过期定时任务每分钟执行
□ 会员降级定时任务每日执行
□ 全链路 trace_id 透传完成
□ 日志保留策略已配置
□ 压力测试:1000并发消耗请求,余额一致
□ 故障演练:主库宕机切换、快照重建恢复
□ 幂等测试:相同请求连续发送3次,只生效1次
总结:以上设计以 事件溯源 为核心架构,通过 不可变事件表 + 快照表 的双表校验机制保证账目绝对准确;通过 三层幂等防护(Redis → DB唯一索引 → 业务查重)保证操作幂等;通过 每日自动对账 + 快照重建 保证可回溯和可修复;通过 悲观锁 + 分布式锁 + 数据库事务 保证并发安全。这套方案可以达到金融级账目系统的可靠性标准。