refactor: 流量系统重构 — 增量累加算法 + 日粒度缓冲 + 旧详单清理
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m16s
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m16s
核心改造: - 增量算法:流量计算从覆盖式改为增量累加(gateway - lastReading),支持上游运营商重置检测 - 日流量缓冲:insertDataUsageRecord 改为 Redis INCRBYFLOAT,每日凌晨落盘到 tb_card_daily_usage - 运营商:新增 data_reset_day 字段(联通=27,其余=1) - IoT卡:新增 last_gateway_reading_mb 字段存储上次网关读数 - 查询层:新建 TrafficQueryService 合并 Redis(今日)+ DB(历史)数据源 - 清理:删除 DataUsageRecord model/store,移除 polling_handler 旧引用 迁移:000094-000097(carrier字段、iot_card字段、数据初始化、日流量表)
This commit is contained in:
151
internal/task/daily_traffic_flush.go
Normal file
151
internal/task/daily_traffic_flush.go
Normal file
@@ -0,0 +1,151 @@
|
||||
package task
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/redis/go-redis/v9"
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/break/junhong_cmp_fiber/internal/model"
|
||||
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
||||
)
|
||||
|
||||
// DailyTrafficFlushHandler 每日流量落盘任务处理器
|
||||
// 凌晨 2 点分批 SCAN 昨日 Redis key → UPSERT 到 tb_card_daily_usage → 删 key
|
||||
type DailyTrafficFlushHandler struct {
|
||||
redis *redis.Client
|
||||
cardDailyUsageStore *postgres.CardDailyUsageStore
|
||||
logger *zap.Logger
|
||||
}
|
||||
|
||||
func NewDailyTrafficFlushHandler(
|
||||
redisClient *redis.Client,
|
||||
cardDailyUsageStore *postgres.CardDailyUsageStore,
|
||||
logger *zap.Logger,
|
||||
) *DailyTrafficFlushHandler {
|
||||
return &DailyTrafficFlushHandler{
|
||||
redis: redisClient,
|
||||
cardDailyUsageStore: cardDailyUsageStore,
|
||||
logger: logger,
|
||||
}
|
||||
}
|
||||
|
||||
func (h *DailyTrafficFlushHandler) HandleDailyTrafficFlush(ctx context.Context, _ *asynq.Task) error {
|
||||
startTime := time.Now()
|
||||
|
||||
yesterday := time.Now().AddDate(0, 0, -1).Format("2006-01-02")
|
||||
pattern := fmt.Sprintf("traffic:daily:*:%s", yesterday)
|
||||
|
||||
h.logger.Info("开始每日流量落盘",
|
||||
zap.String("date", yesterday),
|
||||
zap.String("pattern", pattern))
|
||||
|
||||
// 分批 SCAN 收集所有昨日 key
|
||||
var allKeys []string
|
||||
var cursor uint64
|
||||
for {
|
||||
keys, nextCursor, err := h.redis.Scan(ctx, cursor, pattern, 500).Result()
|
||||
if err != nil {
|
||||
h.logger.Error("SCAN Redis key 失败", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
allKeys = append(allKeys, keys...)
|
||||
cursor = nextCursor
|
||||
if cursor == 0 {
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
if len(allKeys) == 0 {
|
||||
h.logger.Info("无需落盘,昨日无流量数据", zap.String("date", yesterday))
|
||||
return nil
|
||||
}
|
||||
|
||||
h.logger.Info("收集到待落盘 key",
|
||||
zap.String("date", yesterday),
|
||||
zap.Int("count", len(allKeys)))
|
||||
|
||||
yesterdayDate, _ := time.Parse("2006-01-02", yesterday)
|
||||
|
||||
// 分批处理:每批 200 条
|
||||
flushedCount := 0
|
||||
for i := 0; i < len(allKeys); i += 200 {
|
||||
end := i + 200
|
||||
if end > len(allKeys) {
|
||||
end = len(allKeys)
|
||||
}
|
||||
batch := allKeys[i:end]
|
||||
|
||||
var records []*model.CardDailyUsage
|
||||
var flushedKeys []string
|
||||
|
||||
for _, key := range batch {
|
||||
cardID, err := parseCardIDFromKey(key)
|
||||
if err != nil {
|
||||
h.logger.Warn("解析 Redis key 失败", zap.String("key", key), zap.Error(err))
|
||||
continue
|
||||
}
|
||||
|
||||
val, err := h.redis.Get(ctx, key).Float64()
|
||||
if err != nil {
|
||||
h.logger.Warn("读取 Redis 值失败", zap.String("key", key), zap.Error(err))
|
||||
continue
|
||||
}
|
||||
|
||||
records = append(records, &model.CardDailyUsage{
|
||||
IotCardID: cardID,
|
||||
Date: yesterdayDate,
|
||||
UsageMB: val,
|
||||
})
|
||||
flushedKeys = append(flushedKeys, key)
|
||||
}
|
||||
|
||||
// 批量 UPSERT 到 DB
|
||||
if len(records) > 0 {
|
||||
if err := h.cardDailyUsageStore.UpsertBatch(ctx, records); err != nil {
|
||||
h.logger.Error("批量 UPSERT 失败",
|
||||
zap.Int("batch_size", len(records)),
|
||||
zap.Error(err))
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
// 批量删除已落盘的 Redis key
|
||||
if len(flushedKeys) > 0 {
|
||||
pipe := h.redis.Pipeline()
|
||||
for _, key := range flushedKeys {
|
||||
pipe.Del(ctx, key)
|
||||
}
|
||||
if _, err := pipe.Exec(ctx); err != nil {
|
||||
h.logger.Warn("批量删除 Redis key 失败", zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
flushedCount += len(records)
|
||||
}
|
||||
|
||||
h.logger.Info("每日流量落盘完成",
|
||||
zap.String("date", yesterday),
|
||||
zap.Int("flushed_count", flushedCount),
|
||||
zap.Duration("duration", time.Since(startTime)))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// parseCardIDFromKey 从 Redis key "traffic:daily:{cardID}:{date}" 中解析 cardID
|
||||
func parseCardIDFromKey(key string) (uint, error) {
|
||||
parts := strings.Split(key, ":")
|
||||
if len(parts) != 4 {
|
||||
return 0, fmt.Errorf("invalid key format: %s", key)
|
||||
}
|
||||
id, err := strconv.ParseUint(parts[2], 10, 64)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return uint(id), nil
|
||||
}
|
||||
Reference in New Issue
Block a user