跳转至

家纺AI生图平台 · 积分与会员系统技术设计

文档定位:目标架构(事件溯源 + 快照 + 幂等 + 对账)。现行实现为「会员方案 → 单一 usage_quotas 扣点」,见 membership-agent-sales-guide.md;本文描述金融级账本目标,不代表已全部落地。

目标:确保账目 准确性 · 持久性 · 健壮性 · 可回溯 · 幂等 · 日查,达到金融级可靠性


目录

  1. 设计原则
  2. 数据模型设计
  3. 幂等性设计
  4. 事务与并发控制
  5. 审计与可回溯
  6. 日终对账体系
  7. 持久性与容灾
  8. 健壮性保障
  9. 关键业务流程
  10. 监控与告警
  11. 技术选型建议

一、设计原则

┌──────────────────────────────────────────────────┐
│                金融级账目系统五原则                  │
├──────────────────────────────────────────────────┤
│  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唯一索引 → 业务查重)保证操作幂等;通过 每日自动对账 + 快照重建 保证可回溯和可修复;通过 悲观锁 + 分布式锁 + 数据库事务 保证并发安全。这套方案可以达到金融级账目系统的可靠性标准。