From 74f43188388f87217aec840a1ea10d11ce5b1865 Mon Sep 17 00:00:00 2001 From: eafonyang Date: Wed, 10 Jun 2026 19:28:30 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96=E7=BC=93=E5=AD=98=E7=BB=93?= =?UTF-8?q?=E6=9E=84=EF=BC=8C=E4=BC=98=E5=8C=96=E7=BC=93=E5=AD=98=E6=9C=BA?= =?UTF-8?q?=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/cache/curve_cache.go | 63 ++++++++++++++- internal/handler/curve.go | 139 ++++++++++++++++++++++++++-------- 2 files changed, 169 insertions(+), 33 deletions(-) diff --git a/internal/cache/curve_cache.go b/internal/cache/curve_cache.go index 231618b..aa20e90 100644 --- a/internal/cache/curve_cache.go +++ b/internal/cache/curve_cache.go @@ -2,6 +2,7 @@ package cache import ( "context" + "encoding/json" "fmt" "time" @@ -9,7 +10,7 @@ import ( ) const ( - curveLockTTL = 10 * time.Second + curveLockTTL = 10 * time.Second curveLockRetryDelay = 100 * time.Millisecond curveLockMaxRetries = 50 ) @@ -42,6 +43,64 @@ func (c *CurveCache) Set(ctx context.Context, brand, name, target, data string) return c.rdb.HSet(ctx, key, target, data).Err() } +// GetFR 直接获取独立存储的 fr 数据(用于 modelCurve 接口) +func (c *CurveCache) GetFR(ctx context.Context, brand, name string) (string, error) { + key := brand + " " + name + val, err := c.rdb.HGet(ctx, key, "__fr").Result() + if err == redis.Nil { + return "", nil + } + if err != nil { + return "", fmt.Errorf("hget curve fr: %w", err) + } + return val, nil +} + +// GetWithFR 获取 target 缓存数据并合并独立存储的 fr 数据 +// 兼容旧缓存:若 target 数据中已有 fr 字段则直接返回 +func (c *CurveCache) GetWithFR(ctx context.Context, brand, name, target string) (string, error) { + key := brand + " " + name + + targetData, err := c.rdb.HGet(ctx, key, target).Result() + if err == redis.Nil { + return "", nil + } + if err != nil { + return "", fmt.Errorf("hget curve target: %w", err) + } + if targetData == "" { + return "", nil + } + + // 旧缓存兼容:已有 fr 字段无需合并 + var check map[string]json.RawMessage + if err := json.Unmarshal([]byte(targetData), &check); err == nil { + if _, hasFR := check["fr"]; hasFR { + return targetData, nil + } + } + + // 新缓存:从 __fr 字段读取 fr 数据并合并 + frData, err := c.rdb.HGet(ctx, key, "__fr").Result() + if err == redis.Nil || frData == "" { + return targetData, nil + } + if err != nil { + return targetData, nil + } + + var result map[string]json.RawMessage + if err := json.Unmarshal([]byte(targetData), &result); err != nil { + return targetData, nil + } + result["fr"] = json.RawMessage(frData) + merged, err := json.Marshal(result) + if err != nil { + return targetData, nil + } + return string(merged), nil +} + // AcquireLock 获取分布式锁(SETNX),防止并发请求同一个曲线数据 func (c *CurveCache) AcquireLock(ctx context.Context, brand, name, target string) (bool, error) { lockKey := brand + " " + name + ":" + target + ":lock" @@ -85,4 +144,4 @@ func (c *CurveCache) GetWithLock(ctx context.Context, brand, name, target string // 未获取锁,等待后重试 return "", false, nil -} \ No newline at end of file +} diff --git a/internal/handler/curve.go b/internal/handler/curve.go index 4d78640..f4cfc01 100644 --- a/internal/handler/curve.go +++ b/internal/handler/curve.go @@ -55,6 +55,26 @@ func (h *CurveHandler) ModelCurve(c *gin.Context) { ctx := c.Request.Context() + // 快速路径:直接读取独立存储的 fr 数据(缓存命中时只需一次 Redis 调用) + frData, err := h.cache.GetFR(ctx, brand, name) + if err != nil { + h.log.Warn("get fr from cache failed", zap.Error(err)) + } + + if frData != "" { + h.log.Info("model curve fr cache hit", zap.String("brand", brand), zap.String("name", name)) + var fr any + if err := json.Unmarshal([]byte(frData), &fr); err == nil { + h.writeCurveResponse(c, base64Resp, gin.H{ + "code": 200, + "msg": "ok", + "fr": fr, + }) + return + } + } + + // 慢速路径:缓存未命中,走完整流程获取数据(同时会写入 __fr 缓存) result, err := h.getCurvePoint(ctx, brand, name, defaultModelCurveTarget) if err != nil { h.log.Error("get model curve failed", zap.Error(err)) @@ -76,7 +96,21 @@ func (h *CurveHandler) ModelCurve(c *gin.Context) { return } - h.writeCurveResponse(c, base64Resp, resp) + // 只返回 fr,不包含 parametric_eq + fr, hasFR := resp["fr"] + if !hasFR { + h.writeCurveResponse(c, base64Resp, gin.H{ + "code": 0, + "msg": "无曲线数据", + }) + return + } + + h.writeCurveResponse(c, base64Resp, gin.H{ + "code": 200, + "msg": "ok", + "fr": fr, + }) } // GetCurve 获取目标曲线 @@ -151,47 +185,90 @@ func (h *CurveHandler) writeCurveResponse(c *gin.Context, base64Resp bool, data } // getCurvePoint 获取曲线数据:先查缓存,缓存不存在则请求 EQ 接口并缓存结果 +// 缓存结构优化:fr 数据(与 target 无关)单独存储在 __fr 字段,避免每个 target 重复存储 func (h *CurveHandler) getCurvePoint(ctx context.Context, brand, name, target string) (string, error) { - data, acquired, err := h.cache.GetWithLock(ctx, brand, name, target) + // 使用 GetWithFR,命中缓存时自动合并 __fr 数据 + data, err := h.cache.GetWithFR(ctx, brand, name, target) if err != nil { return "", err } - - // 缓存命中 - if data != "" && !acquired { + if data != "" { h.log.Info("curve cache hit", zap.String("brand", brand), zap.String("name", name), zap.String("target", target)) return data, nil } - // 获取了锁,缓存仍然为空,需要请求 EQ 接口 - if acquired { - h.log.Info("curve cache miss, requesting eq api", zap.String("brand", brand), zap.String("name", name), zap.String("target", target)) - result, eqErr := h.getCurvePointFromPEQ(ctx, brand, name, target) - if eqErr != nil { - // 释放锁 - if lockErr := h.cache.ReleaseLock(ctx, brand, name, target); lockErr != nil { - h.log.Warn("release curve lock failed", zap.Error(lockErr)) - } - return "", eqErr - } - // 释放锁 - if lockErr := h.cache.ReleaseLock(ctx, brand, name, target); lockErr != nil { - h.log.Warn("release curve lock failed", zap.Error(lockErr)) - } - - if result != "" { - // 缓存结果 - if cacheErr := h.cache.Set(ctx, brand, name, target, result); cacheErr != nil { - h.log.Warn("curve cache set failed", zap.Error(cacheErr)) - } - return result, nil - } - return "", nil + // 缓存未命中,获取分布式锁 + acquired, lockErr := h.cache.AcquireLock(ctx, brand, name, target) + if lockErr != nil { + return "", lockErr } // 未获取锁(其他请求正在处理),等待后重试 - time.Sleep(100 * time.Millisecond) - return h.getCurvePoint(ctx, brand, name, target) + if !acquired { + time.Sleep(100 * time.Millisecond) + return h.getCurvePoint(ctx, brand, name, target) + } + + // 获取了锁,双重检查缓存 + h.log.Info("curve cache miss, requesting eq api", zap.String("brand", brand), zap.String("name", name), zap.String("target", target)) + data, err = h.cache.GetWithFR(ctx, brand, name, target) + if err != nil { + h.releaseLock(brand, name, target) + return "", err + } + if data != "" { + h.releaseLock(brand, name, target) + return data, nil + } + + // 请求 EQ 接口 + result, eqErr := h.getCurvePointFromPEQ(ctx, brand, name, target) + if eqErr != nil { + h.releaseLock(brand, name, target) + return "", eqErr + } + + if result == "" { + h.releaseLock(brand, name, target) + return "", nil + } + + // 将 fr 数据剥离,单独存储在 __fr 字段,避免每个 target 重复存储 + dataToCache := result + var frToCache string + var parsed map[string]json.RawMessage + if err := json.Unmarshal([]byte(result), &parsed); err == nil { + if fr, ok := parsed["fr"]; ok { + frToCache = string(fr) + delete(parsed, "fr") + if stripped, err := json.Marshal(parsed); err == nil { + dataToCache = string(stripped) + } + } + } + + // 先写缓存,再释放锁(避免其他服务器在缓存写入前抢到锁后重复调用 EQ API) + if cacheErr := h.cache.Set(ctx, brand, name, target, dataToCache); cacheErr != nil { + h.log.Warn("curve cache set failed", zap.Error(cacheErr)) + } + if frToCache != "" { + existingFR, _ := h.cache.Get(ctx, brand, name, "__fr") + if existingFR == "" { + if cacheErr := h.cache.Set(ctx, brand, name, "__fr", frToCache); cacheErr != nil { + h.log.Warn("curve fr cache set failed", zap.Error(cacheErr)) + } + } + } + + h.releaseLock(brand, name, target) + return result, nil +} + +// releaseLock 释放分布式锁并记录警告日志 +func (h *CurveHandler) releaseLock(brand, name, target string) { + if err := h.cache.ReleaseLock(context.Background(), brand, name, target); err != nil { + h.log.Warn("release curve lock failed", zap.Error(err)) + } } // getCurvePointFromPEQ 从 EQ 接口获取曲线数据