「产品矩阵平台」数据与缓存架构

产品矩阵数据与缓存架构:MySQL 数据库分层、读写分离、Redis 分布式缓存、ClickHouse 分析处理与 ES 全文检索的完整数据架构设计。

第七章 数据与缓存架构

数据架构决定平台能跑多远。早期系统可以靠 MySQL 单库支撑,但产品矩阵一旦出现多应用、多租户、内容增长、行为分析和 AI 调用,数据会迅速分层。

建议把数据分成四类:

类型存储特点
交易数据MySQL / PostgreSQL强一致、可审计
缓存数据Redis / Ristretto低延迟、可过期
搜索数据ElasticSearch / OpenSearch倒排索引、全文检索
分析数据ClickHouse / Hive大宽表、聚合查询

7.1 数据库分层架构

基础形态:

graph TD
    A[Application] --> B[Repository]
    B --> C[Primary DB]
    B --> D[Read Replica]
    B --> E[Redis Cache]
    C --> F[Binlog / CDC]
    F --> G[Search Index]
    F --> H[Analytics Store]

读写分离要谨慎。并不是所有读请求都能走从库,例如刚创建订单后立即查询订单详情,如果读到延迟从库,会出现“下单成功但看不到订单”的体验。

推荐策略:

场景读取位置
用户刚写入后的强一致查询主库
后台列表、报表预览从库
内容详情、商品详情缓存优先
BI 聚合ClickHouse

7.2 ORM Hook 实践

GORM Hook 可以增强一致性,但要控制边界。

适合放在 Hook 的逻辑:

能力原因
tenant_id 自动填充横切规则,避免漏写
created_byupdated_by审计字段统一
软删除限制防止误删
更新时间标准化

不适合放在 Hook 的逻辑:

逻辑原因
发送短信外部副作用,不可控
创建订单后扣库存业务流程应显式
复杂权限判断依赖上下文,难测试
调用第三方 API事务中调用风险大

Hook 应该像安全带,而不是方向盘。

7.3 缓存策略

常见缓存模式:

模式写法适用
Cache Aside应用读缓存,未命中读库再写缓存最常用
Write Through写库同时写缓存配置、字典
Write Behind先写缓存,异步落库计数、日志,需谨慎
Refresh Ahead过期前刷新热点数据

Cache Aside 的关键不是“加 Redis”,而是处理缓存击穿、穿透和雪崩。

问题表现防护
击穿热点 key 失效瞬间大量打库单飞锁、提前刷新
穿透查询不存在数据反复打库空值缓存、布隆过滤器
雪崩大量 key 同时过期TTL 加随机抖动

7.4 多级缓存设计

多级缓存适合高读场景,例如配置、权限、商品详情、内容详情。

Local Cache -> Redis -> Database

本地缓存可以用 Ristretto,Redis 用于跨实例共享。配置变更后通过 Pub/Sub 或事件通知各实例失效本地缓存。

缓存 key 设计:

tenant:{tenant_id}:app:{app_id}:content:{content_id}:v{version}

key 中带版本号,可以让发布和回滚更可控。

7.5 数据一致性策略

产品矩阵平台中,强一致只保留给关键交易链路。其他场景优先采用最终一致。

场景一致性要求推荐方案
支付状态强一致 + 对账兜底事务、幂等、主动查询
库存扣减强一致或近强一致锁库、乐观锁
通知发送最终一致事件 + 重试
搜索索引最终一致CDC / Outbox
BI 报表延迟一致批处理

分布式事务不要轻易引入。多数业务可以通过本地事务 + Outbox + 补偿任务解决。

7.6 搜索与分析系统

搜索和分析是两个不同系统:

系统目标查询特点
ElasticSearch找到相关内容全文检索、模糊匹配、排序
ClickHouse统计和分析聚合、分组、时间窗口

不要用 ElasticSearch 做长期 BI,也不要用 ClickHouse 承担实时全文搜索。工具选错,会让维护成本持续上升。

7.7 大数据与埋点管道

埋点事件建议使用统一结构:

{
  "event": "content_viewed",
  "tenant_id": "t_001",
  "app_id": "app_001",
  "user_id": "u_001",
  "anonymous_id": "anon_001",
  "timestamp": "2025-10-16T16:28:36+08:00",
  "properties": {
    "content_id": "c_001",
    "source": "feed"
  }
}

事件进入 Kafka 后,按用途分流:

流向用途
ClickHouse 明细表行为分析
用户画像服务标签更新
推荐服务实时反馈
风控服务异常行为检测

7.8 数据脱敏与隐私保护

敏感数据要先分类:

级别示例处理
公开昵称、公开文章正常展示
内部用户 ID、订单号权限控制
敏感手机号、邮箱、身份证脱敏、加密
高敏支付凭证、密钥独立存储、严格审计

常见脱敏:

手机号:138****5678
邮箱:le***@example.com
身份证:110***********1234

日志中禁止直接打印高敏数据。对调试确实需要的字段,应使用 token 化后的引用 ID。

7.9 数据架构落地清单

检查项标准
表设计多租户字段、索引、审计字段齐全
Repository不泄露 ORM 细节
缓存每类 key 有 TTL、失效策略、命名规范
事件关键数据同步不丢事件
搜索索引构建可重放
分析埋点结构稳定,字段有字典
隐私敏感字段加密、脱敏、审计

7.10 代码实践:Cache-Aside 模板与分布式锁

一、Cache-Aside 防击穿模式

package cache

import (
    "context"
    "encoding/json"
    "fmt"
    "time"
)

type CacheAside[T any] struct {
    redis      RedisClient
    keyPrefix  string
    ttl        time.Duration
    lockTTL    time.Duration
}

func (c *CacheAside[T]) Get(ctx context.Context, key string, loader func() (T, error)) (T, error) {
    var result T
    cacheKey := fmt.Sprintf("%s:%s", c.keyPrefix, key)

    // 1. 先读缓存
    val, err := c.redis.Get(ctx, cacheKey).Result()
    if err == nil {
        _ = json.Unmarshal([]byte(val), &result)
        return result, nil
    }

    // 2. 缓存未命中,获取单飞锁
    lockKey := cacheKey + ":lock"
    lockOk, _ := c.redis.SetNX(ctx, lockKey, "1", c.lockTTL).Result()
    if !lockOk {
        // 其他 goroutine 正在回源,稍等再试
        time.Sleep(50 * time.Millisecond)
        return c.Get(ctx, key, loader)
    }
    defer c.redis.Del(ctx, lockKey)

    // 3. 双重检查
    val, err = c.redis.Get(ctx, cacheKey).Result()
    if err == nil {
        _ = json.Unmarshal([]byte(val), &result)
        return result, nil
    }

    // 4. 回源
    result, err = loader()
    if err != nil {
        return result, err
    }

    // 5. 写入缓存(空值也缓存防穿透,TTL 更短)
    data, _ := json.Marshal(result)
    ttl := c.ttl
    if isZero(result) {
        ttl = 5 * time.Minute // 空值缓存 5 分钟
    }
    _ = c.redis.Set(ctx, cacheKey, data, ttl).Err()
    return result, nil
}

func (c *CacheAside[T]) Invalidate(ctx context.Context, key string) error {
    cacheKey := fmt.Sprintf("%s:%s", c.keyPrefix, key)
    return c.redis.Del(ctx, cacheKey).Err()
}

二、多级缓存(本地 + Redis)

package cache

import (
    "context"
    "time"
    "github.com/dgraph-io/ristretto"
)

type MultiLevelCache[T any] struct {
    local      *ristretto.Cache
    redis      RedisClient
    ttlLocal   time.Duration
    ttlRedis   time.Duration
    lockTTL    time.Duration
    keyPrefix  string
}

func (c *MultiLevelCache[T]) Get(ctx context.Context, key string, loader func() (T, error)) (T, error) {
    var result T
    cacheKey := fmt.Sprintf("%s:%s", c.keyPrefix, key)

    // 1. 先查本地缓存(Ristretto,零 GC,纳秒级)
    if val, ok := c.local.Get(cacheKey); ok {
        return val.(T), nil
    }

    // 2. 再查 Redis
    val, err := c.redis.Get(ctx, cacheKey).Result()
    if err == nil {
        _ = json.Unmarshal([]byte(val), &result)
        c.local.SetWithTTL(cacheKey, result, 1, c.ttlLocal)
        return result, nil
    }

    // 3. 回源(Cache-Aside + singleflight)
    result, err = c.loadWithSingleflight(ctx, cacheKey, loader)
    if err != nil {
        return result, err
    }

    // 4. 回填两级缓存
    data, _ := json.Marshal(result)
    _ = c.redis.Set(ctx, cacheKey, data, c.ttlRedis).Err()
    c.local.SetWithTTL(cacheKey, result, 1, c.ttlLocal)
    return result, nil
}

// 订阅配置变更,通知实例失效本地缓存
func (c *MultiLevelCache[T]) SubscribeInvalidate(channel string) {
    pubsub := c.redis.Subscribe(context.Background(), channel)
    go func() {
        for msg := range pubsub.Channel() {
            c.local.Del(msg.Payload)
        }
    }()
}

三、Redis RedLock 分布式锁

package lock

import (
    "context"
    "fmt"
    "time"
    "github.com/go-redsync/redsync/v4"
    "github.com/go-redsync/redsync/v4/redis/goredis/v9"
)

type DistributedLock struct {
    rs    *redsync.Redsync
    mutex *redsync.Mutex
}

func NewDistributedLock(pool *redis.Pool, resource string, ttl time.Duration) *DistributedLock {
    rs := redsync.New(goredis.NewPool(pool))
    return &DistributedLock{
        rs:    rs,
        mutex: rs.NewMutex(resource, redsync.WithExpiry(ttl)),
    }
}

func (l *DistributedLock) Do(ctx context.Context, fn func() error) error {
    if err := l.mutex.LockContext(ctx); err != nil {
        return fmt.Errorf("lock failed: %w", err)
    }
    defer l.mutex.UnlockContext(ctx)
    return fn()
}

// 使用示例:库存扣减
func (svc *InventoryService) Deduct(ctx context.Context, productID string, qty int) error {
    lock := NewDistributedLock(svc.redisPool, "lock:inventory:"+productID, 5*time.Second)
    return lock.Do(ctx, func() error {
        current, err := svc.GetStock(ctx, productID)
        if err != nil {
            return err
        }
        if current < qty {
            return fmt.Errorf("insufficient stock: %d < %d", current, qty)
        }
        return svc.UpdateStock(ctx, productID, current-qty)
    })
}

四、Outbox 模式确保事件不丢失

package outbox

import (
    "context"
    "encoding/json"
    "gorm.io/gorm"
    "time"
)

type OutboxEvent struct {
    ID        string    `gorm:"primaryKey"`
    EventType string    `gorm:"index"`
    TenantID  string    `gorm:"index"`
    AppID     string    `gorm:"index"`
    Payload   []byte
    CreatedAt time.Time
    SentAt    *time.Time `gorm:"index"`
}

// PublishInTransaction 在业务事务内写 Outbox
type Publisher struct {
    db *gorm.DB
}

func (p *Publisher) PublishInTransaction(ctx context.Context, eventType, tenantID, appID string, payload interface{}) error {
    data, _ := json.Marshal(payload)
    return p.db.WithContext(ctx).Create(&OutboxEvent{
        ID:        uuid.New().String(),
        EventType: eventType,
        TenantID:  tenantID,
        AppID:     appID,
        Payload:   data,
        CreatedAt: time.Now(),
    }).Error
}

// RelayWorker 扫描未投递事件并发送到消息队列
func (p *Publisher) RelayWorker(ctx context.Context, sender func(OutboxEvent) error) {
    ticker := time.NewTicker(5 * time.Second)
    defer ticker.Stop()
    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            var events []OutboxEvent
            p.db.Where("sent_at IS NULL").Order("created_at").Limit(100).Find(&events)
            for _, e := range events {
                if err := sender(e); err != nil {
                    continue // 重试下一批
                }
                now := time.Now()
                p.db.Model(&e).Update("sent_at", now)
            }
        }
    }
}

五、数据脱敏工具函数

package desensitization

import (
    "regexp"
    "strings"
)

func Phone(phone string) string {
    if len(phone) != 11 {
        return phone
    }
    return phone[:3] + "****" + phone[7:]
}

func Email(email string) string {
    parts := strings.Split(email, "@")
    if len(parts) != 2 {
        return email
    }
    local := parts[0]
    if len(local) <= 1 {
        return email
    }
    return local[:1] + strings.Repeat("*", len(local)-1) + "@" + parts[1]
}

func IDCard(id string) string {
    if len(id) != 18 {
        return id
    }
    return id[:6] + "********" + id[14:]
}

// 在 Repository 中封装脱敏
func DesensitizeUser(user *User) {
    user.Phone = Phone(user.Phone)
    user.Email = Email(user.Email)
    user.IDCard = IDCard(user.IDCard)
}

本章小结

本章将平台数据分为交易、缓存、搜索、分析四类,分别对应 MySQL、Redis/Ristretto、ElasticSearch 和 ClickHouse。核心原则包括:读写分离策略按场景区分(强一致走主库)、缓存一致性先于技术选型、搜索与分析各司其职。数据架构落地需要表设计规范、Repository 不泄露 ORM、缓存有 TTL 和命名规范、以及敏感字段加密脱敏。


延伸阅读


关联专题

专题关联内容链接
PostgreSQL数据库分层与多租户Schema/posts/postgresql/
Docker数据存储容器化部署/posts/docker/
GolangGORM Hook与缓存中间件/golang/
Node.js多语言缓存策略复用/posts/nodejs/
CloudflareCDN缓存与边缘数据分发/posts/cloudflare/

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「SaaS」更多文章

  1. 「产品矩阵平台」未来演进方向
  2. 「产品矩阵平台」运维与成本优化
  3. 「产品矩阵平台」安全与合规体系