From 589116b14121352a075c32bddbc722be9c7c2229 Mon Sep 17 00:00:00 2001 From: yiyiis Date: Sat, 23 May 2026 09:23:14 +0800 Subject: [PATCH 1/2] =?UTF-8?q?fix:=20Worker=20=E8=BF=9B=E7=A8=8B=E7=8B=AC?= =?UTF-8?q?=E7=AB=8B=20Channel=20+=20=E8=BF=9B=E7=A8=8B=E5=86=85=E6=8C=87?= =?UTF-8?q?=E6=95=B0=E9=80=80=E9=81=BF=E9=87=8D=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 每个 Worker 消费者使用独立 AMQP Channel,互不影响 - 用 runWorkerWithRetry 包装每个 Worker,断开自动重连 - 用进程内指数退避重试替代失效的 Nack+DLX 重试机制 - 拓扑声明改用临时 Channel,声明完即关闭 --- backend/cmd/worker/main.go | 95 ++++++++++++------- backend/internal/worker/commentworker.go | 27 ++++-- backend/internal/worker/likeworker.go | 26 +++-- backend/internal/worker/notificationworker.go | 26 +++-- backend/internal/worker/popularityworker.go | 27 ++++-- backend/internal/worker/socialworker.go | 27 ++++-- 6 files changed, 154 insertions(+), 74 deletions(-) diff --git a/backend/cmd/worker/main.go b/backend/cmd/worker/main.go index c68e936..6ecdc76 100644 --- a/backend/cmd/worker/main.go +++ b/backend/cmd/worker/main.go @@ -55,6 +55,38 @@ func connectWithRetry(name string, maxRetries int, fn func() error) { log.Fatalf("%s: 超过最大重试次数", name) } +// runWorkerWithRetry 为每个 Worker 创建独立 Channel,断开后自动重连 +func runWorkerWithRetry(ctx context.Context, name string, conn *amqp.Connection, fn func(*amqp.Channel) error) { + for { + select { + case <-ctx.Done(): + return + default: + } + + ch, err := conn.Channel() + if err != nil { + log.Printf("%s: 创建 Channel 失败: %v, 5秒后重试", name, err) + time.Sleep(5 * time.Second) + continue + } + if err := ch.Qos(50, 0, false); err != nil { + log.Printf("%s: QoS 设置失败: %v", name, err) + } + + log.Printf("%s started, consuming", name) + if err := fn(ch); err != nil { + if ctx.Err() != nil { + ch.Close() + return + } + log.Printf("%s: %v, 5秒后重连...", name, err) + } + ch.Close() + time.Sleep(5 * time.Second) + } +} + func main() { // 加载 .env(本地开发) if err := godotenv.Load(); err != nil { @@ -110,42 +142,33 @@ func main() { return err }) defer conn.Close() - // 创建 RabbitMQ 通道 - ch, err := conn.Channel() + + // 用临时 Channel 声明拓扑(持久化队列,声明一次即可) + topoCh, err := conn.Channel() if err != nil { - log.Fatalf("Failed to open rabbitmq channel: %v", err) + log.Fatalf("Failed to open topology channel: %v", err) } - defer ch.Close() - // 声明 Social 交换机和队列 - if err := declareSocialTopology(ch); err != nil { + if err := declareSocialTopology(topoCh); err != nil { log.Fatalf("Failed to declare social topology: %v", err) } - if err := declareLikeTopology(ch); err != nil { + if err := declareLikeTopology(topoCh); err != nil { log.Fatalf("Failed to declare like topology: %v", err) } - if err := declareCommentTopology(ch); err != nil { + if err := declareCommentTopology(topoCh); err != nil { log.Fatalf("Failed to declare comment topology: %v", err) } if cache != nil { - if err := declarePopularityTopology(ch); err != nil { + if err := declarePopularityTopology(topoCh); err != nil { log.Fatalf("Failed to declare popularity topology: %v", err) } } - if err := ch.Qos(50, 0, false); err != nil { - log.Fatalf("Failed to set qos: %v", err) - } + topoCh.Close() - repo := social.NewSocialRepository(sqlDB) - socialWorker := worker.NewSocialWorker(ch, repo, socialQueue) + // 准备 repo + socialRepo := social.NewSocialRepository(sqlDB) videoRepo := video.NewVideoRepository(sqlDB) likeRepo := video.NewLikeRepository(sqlDB) commentRepo := video.NewCommentRepository(sqlDB) - likeWorker := worker.NewLikeWorker(ch, likeRepo, videoRepo, likeQueue) - commentWorker := worker.NewCommentWorker(ch, commentRepo, videoRepo, commentQueue) - var popularityWorker *worker.PopularityWorker - if cache != nil { - popularityWorker = worker.NewPopularityWorker(ch, cache, popularityQueue) - } ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() @@ -162,22 +185,26 @@ func main() { defer pprofServer.Close() } - errCh := make(chan error, 4) - log.Printf("Worker started, consuming queue=%s", socialQueue) - go func() { errCh <- socialWorker.Run(ctx) }() - log.Printf("Worker started, consuming queue=%s", likeQueue) - go func() { errCh <- likeWorker.Run(ctx) }() - log.Printf("Worker started, consuming queue=%s", commentQueue) - go func() { errCh <- commentWorker.Run(ctx) }() - if popularityWorker != nil { - log.Printf("Worker started, consuming queue=%s", popularityQueue) - go func() { errCh <- popularityWorker.Run(ctx) }() + // 每个 Worker 独立 Channel + 自动重连 + go runWorkerWithRetry(ctx, "SocialWorker", conn, func(ch *amqp.Channel) error { + return worker.NewSocialWorker(ch, socialRepo, socialQueue).Run(ctx) + }) + go runWorkerWithRetry(ctx, "LikeWorker", conn, func(ch *amqp.Channel) error { + return worker.NewLikeWorker(ch, likeRepo, videoRepo, likeQueue).Run(ctx) + }) + go runWorkerWithRetry(ctx, "CommentWorker", conn, func(ch *amqp.Channel) error { + return worker.NewCommentWorker(ch, commentRepo, videoRepo, commentQueue).Run(ctx) + }) + if cache != nil { + go runWorkerWithRetry(ctx, "PopularityWorker", conn, func(ch *amqp.Channel) error { + return worker.NewPopularityWorker(ch, cache, popularityQueue).Run(ctx) + }) } - err = <-errCh - if err != nil && err != context.Canceled { - log.Fatalf("Worker stopped: %v", err) - } + // 等待退出信号 + <-ctx.Done() + log.Printf("Worker shutting down...") + time.Sleep(2 * time.Second) // 等待正在处理的消息完成 log.Printf("Worker stopped") } diff --git a/backend/internal/worker/commentworker.go b/backend/internal/worker/commentworker.go index dbfdd3d..1ac8c73 100644 --- a/backend/internal/worker/commentworker.go +++ b/backend/internal/worker/commentworker.go @@ -8,6 +8,7 @@ import ( "feedsystem_video_go/internal/video" "log" "strings" + "time" amqp "github.com/rabbitmq/amqp091-go" ) @@ -58,18 +59,28 @@ func (w *CommentWorker) Run(ctx context.Context) error { } func (w *CommentWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { - if err := w.process(ctx, d.Body); err != nil { - retryCount := rabbitmq.GetRetryCount(d) - if retryCount >= rabbitmq.MaxRetryCount { - log.Printf("comment worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) - _ = d.Ack(false) + const maxRetries = 3 + for i := 0; i <= maxRetries; i++ { + select { + case <-ctx.Done(): + _ = d.Nack(false, true) return + default: } - log.Printf("comment worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) - _ = d.Nack(false, true) + if err := w.process(ctx, d.Body); err != nil { + if i >= maxRetries { + log.Printf("comment worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err) + _ = d.Ack(false) + return + } + wait := time.Duration(1<= rabbitmq.MaxRetryCount { - log.Printf("like worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) - _ = d.Ack(false) + const maxRetries = 3 + for i := 0; i <= maxRetries; i++ { + select { + case <-ctx.Done(): + _ = d.Nack(false, true) return + default: } - log.Printf("like worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) - _ = d.Nack(false, true) + if err := w.process(ctx, d.Body); err != nil { + if i >= maxRetries { + log.Printf("like worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err) + _ = d.Ack(false) + return + } + wait := time.Duration(1<= rabbitmq.MaxRetryCount { - log.Printf("notification worker: max retries, dropping: %v", err) - _ = d.Ack(false) + const maxRetries = 3 + for i := 0; i <= maxRetries; i++ { + select { + case <-ctx.Done(): + _ = d.Nack(false, true) return + default: } - log.Printf("notification worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) - _ = d.Nack(false, true) + if err := w.process(ctx, d); err != nil { + if i >= maxRetries { + log.Printf("notification worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err) + _ = d.Ack(false) + return + } + wait := time.Duration(1<= rabbitmq.MaxRetryCount { - log.Printf("popularity worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) - _ = d.Ack(false) + const maxRetries = 3 + for i := 0; i <= maxRetries; i++ { + select { + case <-ctx.Done(): + _ = d.Nack(false, true) return + default: } - log.Printf("popularity worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) - _ = d.Nack(false, true) + if err := w.process(ctx, d.Body); err != nil { + if i >= maxRetries { + log.Printf("popularity worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err) + _ = d.Ack(false) + return + } + wait := time.Duration(1<= rabbitmq.MaxRetryCount { - log.Printf("social worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) - _ = d.Ack(false) + const maxRetries = 3 + for i := 0; i <= maxRetries; i++ { + select { + case <-ctx.Done(): + _ = d.Nack(false, true) return + default: } - log.Printf("social worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) - _ = d.Nack(false, true) + if err := w.process(ctx, d.Body); err != nil { + if i >= maxRetries { + log.Printf("social worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err) + _ = d.Ack(false) + return + } + wait := time.Duration(1< Date: Sat, 23 May 2026 09:23:47 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix:=20Notification=20=E9=87=8D=E8=BF=9E=20?= =?UTF-8?q?+=20Social=20=E4=BF=AE=E5=A4=8D=20+=20Following=20=E7=BC=93?= =?UTF-8?q?=E5=AD=98=E5=A4=B1=E6=95=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Notification Worker 加重连循环,Channel 断开不再静默死亡 - Social: 先写 DB 再发 MQ,修复 DB 失败产生幽灵通知的问题 - Social: MQ 发布失败记日志,不再静默忽略 - 关注/取关后失效该用户的 Following Feed 缓存(24h TTL) - Redis Client 新增 DelByPattern 按模式批量删除缓存 --- backend/internal/http/router.go | 58 ++++++++-------------- backend/internal/middleware/redis/cache.go | 11 ++++ backend/internal/social/service.go | 53 +++++++++++++++++--- 3 files changed, 76 insertions(+), 46 deletions(-) diff --git a/backend/internal/http/router.go b/backend/internal/http/router.go index c619935..5eeaf21 100644 --- a/backend/internal/http/router.go +++ b/backend/internal/http/router.go @@ -127,7 +127,7 @@ func SetRouter(db *gorm.DB, cache *rediscache.Client, rmq *rabbitmq.RabbitMQ) *g socialMQ = nil } socialRepository := social.NewSocialRepository(db) - socialService := social.NewSocialService(socialRepository, accountRepository, socialMQ) + socialService := social.NewSocialService(socialRepository, accountRepository, socialMQ, cache) socialHandler := social.NewSocialHandler(socialService) socialGroup := r.Group("/social") protectedSocialGroup := socialGroup.Group("") @@ -227,43 +227,25 @@ func SetRouter(db *gorm.DB, cache *rediscache.Client, rmq *rabbitmq.RabbitMQ) *g if rmq != nil { hub := sseHub ctx := context.Background() - // consume from like queue - go func() { - ch, err := rmq.NewChannel() - if err != nil { - log.Printf("notification-like channel: %v", err) - return - } - defer ch.Close() - w := worker.NewNotificationWorker(ch, db, "notification.like", hub) - if err := w.Run(ctx); err != nil { - log.Printf("notification-like worker: %v", err) - } - }() - go func() { - ch, err := rmq.NewChannel() - if err != nil { - log.Printf("notification-comment channel: %v", err) - return - } - defer ch.Close() - w := worker.NewNotificationWorker(ch, db, "notification.comment", hub) - if err := w.Run(ctx); err != nil { - log.Printf("notification-comment worker: %v", err) - } - }() - go func() { - ch, err := rmq.NewChannel() - if err != nil { - log.Printf("notification-social channel: %v", err) - return - } - defer ch.Close() - w := worker.NewNotificationWorker(ch, db, "notification.social", hub) - if err := w.Run(ctx); err != nil { - log.Printf("notification-social worker: %v", err) - } - }() + // 每个 notification worker 独立 Channel + 自动重连 + for _, q := range []string{"notification.like", "notification.comment", "notification.social"} { + go func(queue string) { + for { + ch, err := rmq.NewChannel() + if err != nil { + log.Printf("notification-%s: 创建 Channel 失败: %v, 5秒后重试", queue, err) + time.Sleep(5 * time.Second) + continue + } + w := worker.NewNotificationWorker(ch, db, queue, hub) + if err := w.Run(ctx); err != nil { + log.Printf("notification-%s: %v, 5秒后重连...", queue, err) + } + ch.Close() + time.Sleep(5 * time.Second) + } + }(q) + } } else { log.Printf("Notification SSE disabled (MQ not available)") } diff --git a/backend/internal/middleware/redis/cache.go b/backend/internal/middleware/redis/cache.go index b3e328c..e0369e3 100644 --- a/backend/internal/middleware/redis/cache.go +++ b/backend/internal/middleware/redis/cache.go @@ -27,6 +27,17 @@ func (c *Client) Del(ctx context.Context, key string) error { return c.rdb.Del(ctx, key).Err() } +func (c *Client) DelByPattern(ctx context.Context, pattern string) error { + if c == nil || c.rdb == nil { + return nil + } + iter := c.rdb.Scan(ctx, 0, pattern, 0).Iterator() + for iter.Next(ctx) { + _ = c.rdb.Del(ctx, iter.Val()) + } + return iter.Err() +} + func (c *Client) MGet(cacheCtx context.Context, cacheKeys ...string) ([]interface{}, error) { if c == nil || c.rdb == nil { return nil, errors.New("redis client not initialized") diff --git a/backend/internal/social/service.go b/backend/internal/social/service.go index dfb6147..e7a9b5c 100644 --- a/backend/internal/social/service.go +++ b/backend/internal/social/service.go @@ -5,16 +5,19 @@ import ( "errors" "feedsystem_video_go/internal/account" "feedsystem_video_go/internal/middleware/rabbitmq" + rediscache "feedsystem_video_go/internal/middleware/redis" + "log" ) type SocialService struct { repo *SocialRepository accountrepo *account.AccountRepository socialMQ *rabbitmq.SocialMQ + cache *rediscache.Client } -func NewSocialService(repo *SocialRepository, accountrepo *account.AccountRepository, socialMQ *rabbitmq.SocialMQ) *SocialService { - return &SocialService{repo: repo, accountrepo: accountrepo, socialMQ: socialMQ} +func NewSocialService(repo *SocialRepository, accountrepo *account.AccountRepository, socialMQ *rabbitmq.SocialMQ, cache *rediscache.Client) *SocialService { + return &SocialService{repo: repo, accountrepo: accountrepo, socialMQ: socialMQ, cache: cache} } func (s *SocialService) Follow(ctx context.Context, social *Social) error { @@ -36,10 +39,22 @@ func (s *SocialService) Follow(ctx context.Context, social *Social) error { if isFollowed { return errors.New("already followed") } - if s.socialMQ != nil { - s.socialMQ.Follow(ctx, social.FollowerID, social.VloggerID) + + // 先写 DB,确保数据持久化 + if err := s.repo.Follow(ctx, social); err != nil { + return err } - return s.repo.Follow(ctx, social) + + // DB 成功后,失效该用户的关注列表缓存 + s.invalidateFollowingFeedCache(context.Background(), social.FollowerID) + + // 最后发 MQ(用于通知),失败只记日志不影响业务 + if s.socialMQ != nil { + if err := s.socialMQ.Follow(ctx, social.FollowerID, social.VloggerID); err != nil { + log.Printf("social MQ Follow 发布失败: %v", err) + } + } + return nil } func (s *SocialService) Unfollow(ctx context.Context, social *Social) error { @@ -58,10 +73,32 @@ func (s *SocialService) Unfollow(ctx context.Context, social *Social) error { if !isFollowed { return errors.New("not followed") } - if s.socialMQ != nil { - s.socialMQ.UnFollow(ctx, social.FollowerID, social.VloggerID) + + // 先写 DB + if err := s.repo.Unfollow(ctx, social); err != nil { + return err + } + + // 失效缓存 + s.invalidateFollowingFeedCache(context.Background(), social.FollowerID) + + // 最后发 MQ + if s.socialMQ != nil { + if err := s.socialMQ.UnFollow(ctx, social.FollowerID, social.VloggerID); err != nil { + log.Printf("social MQ UnFollow 发布失败: %v", err) + } + } + return nil +} + +func (s *SocialService) invalidateFollowingFeedCache(ctx context.Context, accountID uint) { + if s.cache == nil { + return + } + pattern := s.cache.Key("feed:listByFollowing:*:accountID=%d:*", accountID) + if err := s.cache.DelByPattern(ctx, pattern); err != nil { + log.Printf("失效 Following 缓存失败: accountID=%d, err=%v", accountID, err) } - return s.repo.Unfollow(ctx, social) } func (s *SocialService) GetAllFollowers(ctx context.Context, VloggerID uint) ([]*account.Account, error) {