fix: 修复轮询系统缓存不一致和可观测性问题
Some checks failed
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Has been cancelled
Some checks failed
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Has been cancelled
- OnCardStatusChanged/Enabled/Disabled 添加 InvalidateCardCache, 解决 Refresh API 更新 DB 后 polling 缓存仍为旧值的 bug - AssetResolveResponse 的 enable_polling/network_status 去掉 omitempty,解决 false/0 时字段从响应中消失的问题 - Scheduler/Initializer/PackageHandler 增加 INFO 级别日志, 可通过日志判断轮询是否工作、处理了哪些卡
This commit is contained in:
@@ -22,11 +22,11 @@ type AssetResolveResponse struct {
|
|||||||
// 状态聚合字段
|
// 状态聚合字段
|
||||||
RealNameStatus int `json:"real_name_status" description:"实名状态:0未实名 1已实名"`
|
RealNameStatus int `json:"real_name_status" description:"实名状态:0未实名 1已实名"`
|
||||||
RealNameAt *time.Time `json:"real_name_at" description:"最近一次完成实名的时间,未实名时为 null"`
|
RealNameAt *time.Time `json:"real_name_at" description:"最近一次完成实名的时间,未实名时为 null"`
|
||||||
CurrentPackage string `json:"current_package" description:"当前套餐名称(无套餐时为空)"`
|
CurrentPackage string `json:"current_package" description:"当前套餐名称(无套餐时为空)"`
|
||||||
PackageTotalMB int64 `json:"package_total_mb" description:"当前套餐总虚流量(MB),已按virtual_ratio换算"`
|
PackageTotalMB int64 `json:"package_total_mb" description:"当前套餐总虚流量(MB),已按virtual_ratio换算"`
|
||||||
PackageUsedMB float64 `json:"package_used_mb" description:"当前已用虚流量(MB),已按virtual_ratio换算"`
|
PackageUsedMB float64 `json:"package_used_mb" description:"当前已用虚流量(MB),已按virtual_ratio换算"`
|
||||||
PackageRemainMB float64 `json:"package_remain_mb" description:"当前套餐剩余虚流量(MB),已按virtual_ratio换算"`
|
PackageRemainMB float64 `json:"package_remain_mb" description:"当前套餐剩余虚流量(MB),已按virtual_ratio换算"`
|
||||||
DeviceProtectStatus string `json:"device_protect_status,omitempty" description:"设备保护期状态:none/stop/start(仅asset_type=device时有效)"`
|
DeviceProtectStatus string `json:"device_protect_status,omitempty" description:"设备保护期状态:none/stop/start(仅asset_type=device时有效)"`
|
||||||
// 绑定关系字段
|
// 绑定关系字段
|
||||||
ICCID string `json:"iccid,omitempty" description:"卡ICCID(asset_type=card时有效)"`
|
ICCID string `json:"iccid,omitempty" description:"卡ICCID(asset_type=card时有效)"`
|
||||||
BoundDeviceID *uint `json:"bound_device_id,omitempty" description:"绑定的设备ID(asset_type=card时有效)"`
|
BoundDeviceID *uint `json:"bound_device_id,omitempty" description:"绑定的设备ID(asset_type=card时有效)"`
|
||||||
@@ -57,8 +57,8 @@ type AssetResolveResponse struct {
|
|||||||
CardCategory string `json:"card_category,omitempty" description:"卡业务类型"`
|
CardCategory string `json:"card_category,omitempty" description:"卡业务类型"`
|
||||||
Supplier string `json:"supplier,omitempty" description:"供应商"`
|
Supplier string `json:"supplier,omitempty" description:"供应商"`
|
||||||
ActivationStatus int `json:"activation_status,omitempty" description:"激活状态"`
|
ActivationStatus int `json:"activation_status,omitempty" description:"激活状态"`
|
||||||
EnablePolling bool `json:"enable_polling,omitempty" description:"是否参与轮询"`
|
EnablePolling bool `json:"enable_polling" description:"是否参与轮询"`
|
||||||
NetworkStatus int `json:"network_status,omitempty" description:"网络状态:0停机 1开机(asset_type=card时有效)"`
|
NetworkStatus int `json:"network_status" description:"网络状态:0停机 1开机(asset_type=card时有效)"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// BoundCardInfo 设备绑定的卡信息
|
// BoundCardInfo 设备绑定的卡信息
|
||||||
@@ -88,25 +88,25 @@ type AssetRealtimeStatusResponse struct {
|
|||||||
|
|
||||||
// AssetPackageResponse 资产套餐信息
|
// AssetPackageResponse 资产套餐信息
|
||||||
type AssetPackageResponse struct {
|
type AssetPackageResponse struct {
|
||||||
PackageUsageID uint `json:"package_usage_id" description:"套餐使用记录ID"`
|
PackageUsageID uint `json:"package_usage_id" description:"套餐使用记录ID"`
|
||||||
PackageID uint `json:"package_id" description:"套餐ID"`
|
PackageID uint `json:"package_id" description:"套餐ID"`
|
||||||
PackageName string `json:"package_name" description:"套餐名称"`
|
PackageName string `json:"package_name" description:"套餐名称"`
|
||||||
PackageType string `json:"package_type" description:"套餐类型:formal/addon"`
|
PackageType string `json:"package_type" description:"套餐类型:formal/addon"`
|
||||||
UsageType string `json:"usage_type" description:"使用类型:single_card/device"`
|
UsageType string `json:"usage_type" description:"使用类型:single_card/device"`
|
||||||
Status int `json:"status" description:"状态:0待生效 1生效中 2已用完 3已过期 4已失效"`
|
Status int `json:"status" description:"状态:0待生效 1生效中 2已用完 3已过期 4已失效"`
|
||||||
StatusName string `json:"status_name" description:"状态名称"`
|
StatusName string `json:"status_name" description:"状态名称"`
|
||||||
DataLimitMB int64 `json:"data_limit_mb" description:"套餐真流量总量(MB)"`
|
DataLimitMB int64 `json:"data_limit_mb" description:"套餐真流量总量(MB)"`
|
||||||
VirtualLimitMB int64 `json:"virtual_limit_mb" description:"套餐虚流量总量(MB),按virtual_ratio换算"`
|
VirtualLimitMB int64 `json:"virtual_limit_mb" description:"套餐虚流量总量(MB),按virtual_ratio换算"`
|
||||||
DataUsageMB int64 `json:"data_usage_mb" description:"已用真流量(MB)"`
|
DataUsageMB int64 `json:"data_usage_mb" description:"已用真流量(MB)"`
|
||||||
VirtualUsedMB float64 `json:"virtual_used_mb" description:"已用虚流量(MB),按virtual_ratio换算"`
|
VirtualUsedMB float64 `json:"virtual_used_mb" description:"已用虚流量(MB),按virtual_ratio换算"`
|
||||||
VirtualRemainMB float64 `json:"virtual_remain_mb" description:"剩余虚流量(MB),按virtual_ratio换算"`
|
VirtualRemainMB float64 `json:"virtual_remain_mb" description:"剩余虚流量(MB),按virtual_ratio换算"`
|
||||||
VirtualRatio float64 `json:"virtual_ratio" description:"虚流量比例(real/virtual)"`
|
VirtualRatio float64 `json:"virtual_ratio" description:"虚流量比例(real/virtual)"`
|
||||||
EnableVirtualData bool `json:"enable_virtual_data" description:"是否启用虚流量"`
|
EnableVirtualData bool `json:"enable_virtual_data" description:"是否启用虚流量"`
|
||||||
ActivatedAt *time.Time `json:"activated_at,omitempty" description:"激活时间(待生效套餐为空)"`
|
ActivatedAt *time.Time `json:"activated_at,omitempty" description:"激活时间(待生效套餐为空)"`
|
||||||
ExpiresAt *time.Time `json:"expires_at,omitempty" description:"到期时间(待生效套餐为空)"`
|
ExpiresAt *time.Time `json:"expires_at,omitempty" description:"到期时间(待生效套餐为空)"`
|
||||||
MasterUsageID *uint `json:"master_usage_id,omitempty" description:"主套餐ID(加油包时有值)"`
|
MasterUsageID *uint `json:"master_usage_id,omitempty" description:"主套餐ID(加油包时有值)"`
|
||||||
Priority int `json:"priority" description:"优先级"`
|
Priority int `json:"priority" description:"优先级"`
|
||||||
CreatedAt time.Time `json:"created_at" description:"创建时间"`
|
CreatedAt time.Time `json:"created_at" description:"创建时间"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// AssetPackagesResult 套餐列表分页结果
|
// AssetPackagesResult 套餐列表分页结果
|
||||||
|
|||||||
@@ -199,6 +199,8 @@ func (p *PollingInitializer) initBatch(ctx context.Context, cards []*model.IotCa
|
|||||||
cardCacheTTL := 7 * 24 * time.Hour
|
cardCacheTTL := 7 * 24 * time.Hour
|
||||||
pipe := p.redis.Pipeline()
|
pipe := p.redis.Pipeline()
|
||||||
cmdCount := 0
|
cmdCount := 0
|
||||||
|
enqueuedCards := 0
|
||||||
|
skippedCards := 0
|
||||||
|
|
||||||
flushPipe := func() {
|
flushPipe := func() {
|
||||||
if cmdCount == 0 {
|
if cmdCount == 0 {
|
||||||
@@ -214,8 +216,10 @@ func (p *PollingInitializer) initBatch(ctx context.Context, cards []*model.IotCa
|
|||||||
for _, card := range cards {
|
for _, card := range cards {
|
||||||
cfg := p.configMgr.MatchConfig(card)
|
cfg := p.configMgr.MatchConfig(card)
|
||||||
if cfg == nil {
|
if cfg == nil {
|
||||||
|
skippedCards++
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
enqueuedCards++
|
||||||
|
|
||||||
shardID := int(card.ID) % p.queueMgr.shardCount
|
shardID := int(card.ID) % p.queueMgr.shardCount
|
||||||
cardIDStr := fmt.Sprintf("%d", card.ID)
|
cardIDStr := fmt.Sprintf("%d", card.ID)
|
||||||
@@ -273,6 +277,10 @@ func (p *PollingInitializer) initBatch(ctx context.Context, cards []*model.IotCa
|
|||||||
}
|
}
|
||||||
|
|
||||||
flushPipe()
|
flushPipe()
|
||||||
|
p.logger.Info("批量初始化完成",
|
||||||
|
zap.Int("total", len(cards)),
|
||||||
|
zap.Int("enqueued", enqueuedCards),
|
||||||
|
zap.Int("skipped_no_config", skippedCards))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -63,11 +63,14 @@ func (s *PollingLifecycleService) OnBatchCardsCreated(ctx context.Context, cards
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// OnCardStatusChanged 卡状态变化后重新匹配配置并更新队列
|
// OnCardStatusChanged 卡状态变化后清理缓存、重新匹配配置并更新队列
|
||||||
func (s *PollingLifecycleService) OnCardStatusChanged(ctx context.Context, cardID uint) {
|
func (s *PollingLifecycleService) OnCardStatusChanged(ctx context.Context, cardID uint) {
|
||||||
if err := s.queueMgr.RemoveFromAllQueues(ctx, cardID); err != nil {
|
if err := s.queueMgr.RemoveFromAllQueues(ctx, cardID); err != nil {
|
||||||
s.logger.Warn("卡状态变化:从队列移除失败", zap.Uint("card_id", cardID), zap.Error(err))
|
s.logger.Warn("卡状态变化:从队列移除失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||||||
}
|
}
|
||||||
|
// 清理轮询卡缓存,确保下次轮询从 DB 读取最新状态
|
||||||
|
s.queueMgr.InvalidateCardCache(ctx, cardID)
|
||||||
|
|
||||||
card, err := s.iotCardStore.GetByID(ctx, cardID)
|
card, err := s.iotCardStore.GetByID(ctx, cardID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.logger.Error("卡状态变化:加载卡信息失败", zap.Uint("card_id", cardID), zap.Error(err))
|
s.logger.Error("卡状态变化:加载卡信息失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||||||
@@ -88,6 +91,7 @@ func (s *PollingLifecycleService) OnCardDeleted(ctx context.Context, cardID uint
|
|||||||
|
|
||||||
// OnCardEnabled 卡启用轮询后初始化
|
// OnCardEnabled 卡启用轮询后初始化
|
||||||
func (s *PollingLifecycleService) OnCardEnabled(ctx context.Context, cardID uint) {
|
func (s *PollingLifecycleService) OnCardEnabled(ctx context.Context, cardID uint) {
|
||||||
|
s.queueMgr.InvalidateCardCache(ctx, cardID)
|
||||||
card, err := s.iotCardStore.GetByID(ctx, cardID)
|
card, err := s.iotCardStore.GetByID(ctx, cardID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.logger.Error("卡启用:加载卡信息失败", zap.Uint("card_id", cardID), zap.Error(err))
|
s.logger.Error("卡启用:加载卡信息失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||||||
@@ -104,6 +108,7 @@ func (s *PollingLifecycleService) OnCardDisabled(ctx context.Context, cardID uin
|
|||||||
if err := s.queueMgr.RemoveFromAllQueues(ctx, cardID); err != nil {
|
if err := s.queueMgr.RemoveFromAllQueues(ctx, cardID); err != nil {
|
||||||
s.logger.Warn("卡禁用:从队列移除失败", zap.Uint("card_id", cardID), zap.Error(err))
|
s.logger.Warn("卡禁用:从队列移除失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||||||
}
|
}
|
||||||
|
s.queueMgr.InvalidateCardCache(ctx, cardID)
|
||||||
}
|
}
|
||||||
|
|
||||||
// shouldEnqueue M3 修复:检查卡或其绑定设备是否允许轮询
|
// shouldEnqueue M3 修复:检查卡或其绑定设备是否允许轮询
|
||||||
|
|||||||
@@ -129,6 +129,14 @@ func (m *PollingQueueManager) OnCardDeleted(ctx context.Context, cardID uint) er
|
|||||||
return m.redis.Del(ctx, cacheKey).Err()
|
return m.redis.Del(ctx, cacheKey).Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// InvalidateCardCache 清理轮询卡信息缓存,强制下次轮询从 DB 重建
|
||||||
|
func (m *PollingQueueManager) InvalidateCardCache(ctx context.Context, cardID uint) {
|
||||||
|
cacheKey := constants.RedisPollingCardInfoKey(cardID)
|
||||||
|
if err := m.redis.Del(ctx, cacheKey).Err(); err != nil {
|
||||||
|
m.logger.Warn("清理轮询卡缓存失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// GetQueueDepth 获取分片队列深度(用于背压检测)
|
// GetQueueDepth 获取分片队列深度(用于背压检测)
|
||||||
func (m *PollingQueueManager) GetQueueDepth(ctx context.Context, shardID int, taskType string) (int64, error) {
|
func (m *PollingQueueManager) GetQueueDepth(ctx context.Context, shardID int, taskType string) (int64, error) {
|
||||||
key := constants.RedisPollingShardQueueKey(shardID, taskType)
|
key := constants.RedisPollingShardQueueKey(shardID, taskType)
|
||||||
|
|||||||
@@ -186,6 +186,9 @@ func (s *Scheduler) processOneShard(ctx context.Context, shardID int) {
|
|||||||
for i, e := range entries {
|
for i, e := range entries {
|
||||||
cardIDs[i] = formatUint(e.CardID)
|
cardIDs[i] = formatUint(e.CardID)
|
||||||
}
|
}
|
||||||
|
s.logger.Info("分片出队",
|
||||||
|
zap.Int("shard_id", shardID), zap.String("task_type", taskType),
|
||||||
|
zap.Int("count", len(entries)))
|
||||||
s.enqueueBatch(ctx, taskType, cardIDs)
|
s.enqueueBatch(ctx, taskType, cardIDs)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -61,13 +61,19 @@ func (h *PollingPackageHandler) Handle(ctx context.Context, t *asynq.Task) error
|
|||||||
}
|
}
|
||||||
|
|
||||||
if h.stopResumeSvc != nil {
|
if h.stopResumeSvc != nil {
|
||||||
// EvaluateAndAct 需要完整新鲜数据(含 DeviceID),从 DB 获取
|
|
||||||
freshCard, loadErr := h.iotCardStore.GetByID(ctx, cardID)
|
freshCard, loadErr := h.iotCardStore.GetByID(ctx, cardID)
|
||||||
if loadErr == nil {
|
if loadErr == nil {
|
||||||
|
h.base.logger.Info("套餐检查:执行停复机评估",
|
||||||
|
zap.Uint("card_id", cardID),
|
||||||
|
zap.Int("network_status", freshCard.NetworkStatus),
|
||||||
|
zap.String("stop_reason", freshCard.StopReason))
|
||||||
if evalErr := h.stopResumeSvc.EvaluateAndAct(ctx, freshCard); evalErr != nil {
|
if evalErr := h.stopResumeSvc.EvaluateAndAct(ctx, freshCard); evalErr != nil {
|
||||||
h.base.logger.Warn("套餐检查后停复机评估失败",
|
h.base.logger.Warn("套餐检查后停复机评估失败",
|
||||||
zap.Uint("card_id", cardID), zap.Error(evalErr))
|
zap.Uint("card_id", cardID), zap.Error(evalErr))
|
||||||
}
|
}
|
||||||
|
} else {
|
||||||
|
h.base.logger.Warn("套餐检查:加载卡信息失败",
|
||||||
|
zap.Uint("card_id", cardID), zap.Error(loadErr))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user