From 9d903ad8e1a519101c8851faf938869c28d5e7b5 Mon Sep 17 00:00:00 2001 From: Sisyphus Date: Sat, 25 Apr 2026 15:54:55 +0800 Subject: [PATCH] =?UTF-8?q?feat(P1):=20MQ=20Worker=20=E6=AD=BB=E4=BF=A1?= =?UTF-8?q?=E9=98=9F=E5=88=97=20=E2=80=94=20=E9=87=8D=E8=AF=95=E4=B8=8A?= =?UTF-8?q?=E9=99=903=E6=AC=A1=E5=90=8E=20Ack=20=E7=A7=BB=E5=85=A5=20DLX?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/internal/middleware/rabbitmq/dlx.go | 53 ++++ .../internal/middleware/rabbitmq/rabbitMQ.go | 239 +++++++-------- backend/internal/worker/commentworker.go | 250 ++++++++-------- backend/internal/worker/likeworker.go | 278 +++++++++--------- backend/internal/worker/popularityworker.go | 164 ++++++----- backend/internal/worker/socialworker.go | 207 ++++++------- 6 files changed, 637 insertions(+), 554 deletions(-) create mode 100644 backend/internal/middleware/rabbitmq/dlx.go diff --git a/backend/internal/middleware/rabbitmq/dlx.go b/backend/internal/middleware/rabbitmq/dlx.go new file mode 100644 index 0000000..82efe5b --- /dev/null +++ b/backend/internal/middleware/rabbitmq/dlx.go @@ -0,0 +1,53 @@ +package rabbitmq + +import ( + "log" + + amqp "github.com/rabbitmq/amqp091-go" +) + +const ( + DLXExchange = "dlx.events" + MaxRetryCount = 3 +) + +// DeclareDLX 声明死信交换机和对应的死信队列 +func DeclareDLX(ch *amqp.Channel, queueName string) error { + if ch == nil { + return nil + } + if err := ch.ExchangeDeclare( + DLXExchange, "topic", true, false, false, false, nil, + ); err != nil { + return err + } + dlxQueue := queueName + ".dlx" + _, err := ch.QueueDeclare( + dlxQueue, true, false, false, false, nil, + ) + if err != nil { + return err + } + if err := ch.QueueBind(dlxQueue, "#", DLXExchange, false, nil); err != nil { + return err + } + log.Printf("DLX ready: exchange=%s queue=%s", DLXExchange, dlxQueue) + return nil +} + +// GetRetryCount 从 AMQP x-death header 中提取当前消息已被重试的次数 +func GetRetryCount(d amqp.Delivery) int { + deaths, ok := d.Headers["x-death"].([]interface{}) + if !ok || len(deaths) == 0 { + return 0 + } + death, ok := deaths[0].(amqp.Table) + if !ok { + return 0 + } + count, ok := death["count"].(int64) + if !ok { + return 0 + } + return int(count) +} diff --git a/backend/internal/middleware/rabbitmq/rabbitMQ.go b/backend/internal/middleware/rabbitmq/rabbitMQ.go index b4a715e..e527ea1 100644 --- a/backend/internal/middleware/rabbitmq/rabbitMQ.go +++ b/backend/internal/middleware/rabbitmq/rabbitMQ.go @@ -1,116 +1,123 @@ -package rabbitmq - -import ( - "context" - "crypto/rand" - "encoding/hex" - "encoding/json" - "errors" - "feedsystem_video_go/internal/config" - "strconv" - "time" - - amqp "github.com/rabbitmq/amqp091-go" -) - -type RabbitMQ struct { - Conn *amqp.Connection - Ch *amqp.Channel -} - -func NewRabbitMQ(cfg *config.RabbitMQConfig) (*RabbitMQ, error) { - if cfg == nil { - return nil, errors.New("rabbitmq config is nil") - } - url := "amqp://" + cfg.Username + ":" + cfg.Password + "@" + cfg.Host + ":" + strconv.Itoa(cfg.Port) + "/" - conn, err := amqp.Dial(url) - if err != nil { - return nil, err - } - ch, err := conn.Channel() - if err != nil { - return nil, err - } - return &RabbitMQ{Conn: conn, Ch: ch}, nil -} - -func (r *RabbitMQ) Close() error { - if r == nil || r.Ch == nil || r.Conn == nil { - return nil - } - if err := r.Ch.Close(); err != nil { - return err - } - if err := r.Conn.Close(); err != nil { - return err - } - return nil -} - -func (r *RabbitMQ) DeclareTopic(exchange string, queue string, bindingKey string) error { - if r == nil || r.Ch == nil { - return errors.New("rabbitmq is not initialized") - } - if exchange == "" || queue == "" || bindingKey == "" { - return errors.New("exchange/queue/bindingKey is required") - } - - if err := r.Ch.ExchangeDeclare( - exchange, - "topic", - true, - false, - false, - false, - nil, - ); err != nil { - return err - } - - q, err := r.Ch.QueueDeclare( - queue, - true, - false, - false, - false, - nil, - ) - if err != nil { - return err - } - - return r.Ch.QueueBind( - q.Name, - bindingKey, - exchange, - false, - nil, - ) -} - -func (r *RabbitMQ) PublishJSON(ctx context.Context, exchange string, routingKey string, payload any) error { - if r == nil || r.Ch == nil { - return errors.New("rabbitmq is not initialized") - } - if exchange == "" || routingKey == "" { - return errors.New("exchange and routingKey are required") - } - b, err := json.Marshal(payload) - if err != nil { - return err - } - return r.Ch.PublishWithContext(ctx, exchange, routingKey, false, false, amqp.Publishing{ - ContentType: "application/json", - DeliveryMode: amqp.Persistent, - Timestamp: time.Now(), - Body: b, - }) -} - -func newEventID(n int) (string, error) { - b := make([]byte, n) - if _, err := rand.Read(b); err != nil { - return "", err - } - return hex.EncodeToString(b), nil -} +package rabbitmq + +import ( + "context" + "crypto/rand" + "encoding/hex" + "encoding/json" + "errors" + "feedsystem_video_go/internal/config" + "log" + "strconv" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +type RabbitMQ struct { + Conn *amqp.Connection + Ch *amqp.Channel +} + +func NewRabbitMQ(cfg *config.RabbitMQConfig) (*RabbitMQ, error) { + if cfg == nil { + return nil, errors.New("rabbitmq config is nil") + } + url := "amqp://" + cfg.Username + ":" + cfg.Password + "@" + cfg.Host + ":" + strconv.Itoa(cfg.Port) + "/" + conn, err := amqp.Dial(url) + if err != nil { + return nil, err + } + ch, err := conn.Channel() + if err != nil { + return nil, err + } + return &RabbitMQ{Conn: conn, Ch: ch}, nil +} + +func (r *RabbitMQ) Close() error { + if r == nil || r.Ch == nil || r.Conn == nil { + return nil + } + if err := r.Ch.Close(); err != nil { + return err + } + if err := r.Conn.Close(); err != nil { + return err + } + return nil +} + +func (r *RabbitMQ) DeclareTopic(exchange string, queue string, bindingKey string) error { + if r == nil || r.Ch == nil { + return errors.New("rabbitmq is not initialized") + } + if exchange == "" || queue == "" || bindingKey == "" { + return errors.New("exchange/queue/bindingKey is required") + } + + if err := r.Ch.ExchangeDeclare( + exchange, + "topic", + true, + false, + false, + false, + nil, + ); err != nil { + return err + } + + q, err := r.Ch.QueueDeclare( + queue, + true, + false, + false, + false, + amqp.Table{"x-dead-letter-exchange": DLXExchange}, + ) + if err != nil { + return err + } + + if err := r.Ch.QueueBind( + q.Name, + bindingKey, + exchange, + false, + nil, + ); err != nil { + return err + } + if err := DeclareDLX(r.Ch, queue); err != nil { + log.Printf("DLX declare failed for %s: %v", queue, err) + } + return nil +} + +func (r *RabbitMQ) PublishJSON(ctx context.Context, exchange string, routingKey string, payload any) error { + if r == nil || r.Ch == nil { + return errors.New("rabbitmq is not initialized") + } + if exchange == "" || routingKey == "" { + return errors.New("exchange and routingKey are required") + } + b, err := json.Marshal(payload) + if err != nil { + return err + } + return r.Ch.PublishWithContext(ctx, exchange, routingKey, false, false, amqp.Publishing{ + ContentType: "application/json", + DeliveryMode: amqp.Persistent, + Timestamp: time.Now(), + Body: b, + }) +} + +func newEventID(n int) (string, error) { + b := make([]byte, n) + if _, err := rand.Read(b); err != nil { + return "", err + } + return hex.EncodeToString(b), nil +} diff --git a/backend/internal/worker/commentworker.go b/backend/internal/worker/commentworker.go index 6b941b3..3fdde7f 100644 --- a/backend/internal/worker/commentworker.go +++ b/backend/internal/worker/commentworker.go @@ -1,122 +1,128 @@ -package worker - -import ( - "context" - "encoding/json" - "errors" - "feedsystem_video_go/internal/middleware/rabbitmq" - "feedsystem_video_go/internal/video" - "log" - "strings" - - amqp "github.com/rabbitmq/amqp091-go" -) - -type CommentWorker struct { - ch *amqp.Channel - comments *video.CommentRepository - videos *video.VideoRepository - queue string -} - -func NewCommentWorker(ch *amqp.Channel, comments *video.CommentRepository, videos *video.VideoRepository, queue string) *CommentWorker { - return &CommentWorker{ch: ch, comments: comments, videos: videos, queue: queue} -} - -func (w *CommentWorker) Run(ctx context.Context) error { - if w == nil || w.ch == nil || w.comments == nil || w.videos == nil { - return errors.New("comment worker is not initialized") - } - if w.queue == "" { - return errors.New("queue is required") - } - - deliveries, err := w.ch.Consume( - w.queue, - "", - false, - false, - false, - false, - nil, - ) - if err != nil { - return err - } - - for { - select { - case <-ctx.Done(): - return ctx.Err() - case d, ok := <-deliveries: - if !ok { - return errors.New("deliveries channel closed") - } - w.handleDelivery(ctx, d) - } - } -} - -func (w *CommentWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { - if err := w.process(ctx, d.Body); err != nil { - log.Printf("comment worker: failed to process message: %v", err) - _ = d.Nack(false, true) - return - } - _ = d.Ack(false) -} - -func (w *CommentWorker) process(ctx context.Context, body []byte) error { - var evt rabbitmq.CommentEvent - if err := json.Unmarshal(body, &evt); err != nil { - return nil - } - switch evt.Action { - case "publish": - return w.applyPublish(ctx, &evt) - case "delete": - return w.applyDelete(ctx, &evt) - default: - return nil - } -} - -func (w *CommentWorker) applyPublish(ctx context.Context, evt *rabbitmq.CommentEvent) error { - if evt == nil || evt.VideoID == 0 || evt.AuthorID == 0 || strings.TrimSpace(evt.Content) == "" { - return nil - } - - ok, err := w.videos.IsExist(ctx, evt.VideoID) - if err != nil { - return err - } - if !ok { - return nil - } - - c := &video.Comment{ - Username: strings.TrimSpace(evt.Username), - VideoID: evt.VideoID, - AuthorID: evt.AuthorID, - Content: strings.TrimSpace(evt.Content), - } - if err := w.comments.CreateComment(ctx, c); err != nil { - return err - } - return w.videos.ChangePopularity(ctx, evt.VideoID, 1) -} - -func (w *CommentWorker) applyDelete(ctx context.Context, evt *rabbitmq.CommentEvent) error { - if evt == nil || evt.CommentID == 0 { - return nil - } - c, err := w.comments.GetByID(ctx, evt.CommentID) - if err != nil { - return err - } - if c == nil { - return nil - } - return w.comments.DeleteComment(ctx, c) -} - +package worker + +import ( + "context" + "encoding/json" + "errors" + "feedsystem_video_go/internal/middleware/rabbitmq" + "feedsystem_video_go/internal/video" + "log" + "strings" + + amqp "github.com/rabbitmq/amqp091-go" +) + +type CommentWorker struct { + ch *amqp.Channel + comments *video.CommentRepository + videos *video.VideoRepository + queue string +} + +func NewCommentWorker(ch *amqp.Channel, comments *video.CommentRepository, videos *video.VideoRepository, queue string) *CommentWorker { + return &CommentWorker{ch: ch, comments: comments, videos: videos, queue: queue} +} + +func (w *CommentWorker) Run(ctx context.Context) error { + if w == nil || w.ch == nil || w.comments == nil || w.videos == nil { + return errors.New("comment worker is not initialized") + } + if w.queue == "" { + return errors.New("queue is required") + } + + deliveries, err := w.ch.Consume( + w.queue, + "", + false, + false, + false, + false, + nil, + ) + if err != nil { + return err + } + + for { + select { + case <-ctx.Done(): + return ctx.Err() + case d, ok := <-deliveries: + if !ok { + return errors.New("deliveries channel closed") + } + w.handleDelivery(ctx, d) + } + } +} + +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) + return + } + log.Printf("comment worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) + _ = d.Nack(false, true) + return + } + _ = d.Ack(false) +} + +func (w *CommentWorker) process(ctx context.Context, body []byte) error { + var evt rabbitmq.CommentEvent + if err := json.Unmarshal(body, &evt); err != nil { + return nil + } + switch evt.Action { + case "publish": + return w.applyPublish(ctx, &evt) + case "delete": + return w.applyDelete(ctx, &evt) + default: + return nil + } +} + +func (w *CommentWorker) applyPublish(ctx context.Context, evt *rabbitmq.CommentEvent) error { + if evt == nil || evt.VideoID == 0 || evt.AuthorID == 0 || strings.TrimSpace(evt.Content) == "" { + return nil + } + + ok, err := w.videos.IsExist(ctx, evt.VideoID) + if err != nil { + return err + } + if !ok { + return nil + } + + c := &video.Comment{ + Username: strings.TrimSpace(evt.Username), + VideoID: evt.VideoID, + AuthorID: evt.AuthorID, + Content: strings.TrimSpace(evt.Content), + } + if err := w.comments.CreateComment(ctx, c); err != nil { + return err + } + return w.videos.ChangePopularity(ctx, evt.VideoID, 1) +} + +func (w *CommentWorker) applyDelete(ctx context.Context, evt *rabbitmq.CommentEvent) error { + if evt == nil || evt.CommentID == 0 { + return nil + } + c, err := w.comments.GetByID(ctx, evt.CommentID) + if err != nil { + return err + } + if c == nil { + return nil + } + return w.comments.DeleteComment(ctx, c) +} + diff --git a/backend/internal/worker/likeworker.go b/backend/internal/worker/likeworker.go index 25a9c6a..ae60369 100644 --- a/backend/internal/worker/likeworker.go +++ b/backend/internal/worker/likeworker.go @@ -1,136 +1,142 @@ -package worker - -import ( - "context" - "encoding/json" - "errors" - "feedsystem_video_go/internal/middleware/rabbitmq" - "feedsystem_video_go/internal/video" - "log" - amqp "github.com/rabbitmq/amqp091-go" - "time" -) - -type LikeWorker struct { - ch *amqp.Channel - likes *video.LikeRepository - videos *video.VideoRepository - queue string -} - -func NewLikeWorker(ch *amqp.Channel, likes *video.LikeRepository, videos *video.VideoRepository, queue string) *LikeWorker { - return &LikeWorker{ch: ch, likes: likes, videos: videos, queue: queue} -} - -func (w *LikeWorker) Run(ctx context.Context) error { - if w == nil || w.ch == nil || w.likes == nil || w.videos == nil { - return errors.New("like worker is not initialized") - } - if w.queue == "" { - return errors.New("queue is required") - } - - deliveries, err := w.ch.Consume( - w.queue, - "", - false, - false, - false, - false, - nil, - ) - if err != nil { - return err - } - - for { - select { - case <-ctx.Done(): - return ctx.Err() - case d, ok := <-deliveries: - if !ok { - return errors.New("deliveries channel closed") - } - w.handleDelivery(ctx, d) - } - } -} - -func (w *LikeWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { - if err := w.process(ctx, d.Body); err != nil { - log.Printf("like worker: failed to process message: %v", err) - _ = d.Nack(false, true) - return - } - _ = d.Ack(false) -} - -func (w *LikeWorker) process(ctx context.Context, body []byte) error { - var evt rabbitmq.LikeEvent - if err := json.Unmarshal(body, &evt); err != nil { - // 解析事件失败,直接丢弃 - return nil - } - if evt.UserID == 0 || evt.VideoID == 0 { - return nil - } - - switch evt.Action { - case "like": - return w.applyLike(ctx, evt.UserID, evt.VideoID) - case "unlike": - return w.applyUnlike(ctx, evt.UserID, evt.VideoID) - default: - return nil - } -} - -func (w *LikeWorker) applyLike(ctx context.Context, userID, videoID uint) error { - ok, err := w.videos.IsExist(ctx, videoID) - if err != nil { - return err - } - if !ok { - return nil - } - - created, err := w.likes.LikeIgnoreDuplicate(ctx, &video.Like{ - VideoID: videoID, - AccountID: userID, - CreatedAt: time.Now(), - }) - if err != nil { - return err - } - if !created { - return nil - } - - if err := w.videos.ChangeLikesCount(ctx, videoID, 1); err != nil { - return err - } - return w.videos.ChangePopularity(ctx, videoID, 1) -} - -func (w *LikeWorker) applyUnlike(ctx context.Context, userID, videoID uint) error { - ok, err := w.videos.IsExist(ctx, videoID) - if err != nil { - return err - } - if !ok { - return nil - } - - deleted, err := w.likes.DeleteByVideoAndAccount(ctx, videoID, userID) - if err != nil { - return err - } - if !deleted { - return nil - } - - if err := w.videos.ChangeLikesCount(ctx, videoID, -1); err != nil { - return err - } - return w.videos.ChangePopularity(ctx, videoID, -1) -} +package worker + +import ( + "context" + "encoding/json" + "errors" + "feedsystem_video_go/internal/middleware/rabbitmq" + "feedsystem_video_go/internal/video" + "log" + amqp "github.com/rabbitmq/amqp091-go" + "time" +) + +type LikeWorker struct { + ch *amqp.Channel + likes *video.LikeRepository + videos *video.VideoRepository + queue string +} + +func NewLikeWorker(ch *amqp.Channel, likes *video.LikeRepository, videos *video.VideoRepository, queue string) *LikeWorker { + return &LikeWorker{ch: ch, likes: likes, videos: videos, queue: queue} +} + +func (w *LikeWorker) Run(ctx context.Context) error { + if w == nil || w.ch == nil || w.likes == nil || w.videos == nil { + return errors.New("like worker is not initialized") + } + if w.queue == "" { + return errors.New("queue is required") + } + + deliveries, err := w.ch.Consume( + w.queue, + "", + false, + false, + false, + false, + nil, + ) + if err != nil { + return err + } + + for { + select { + case <-ctx.Done(): + return ctx.Err() + case d, ok := <-deliveries: + if !ok { + return errors.New("deliveries channel closed") + } + w.handleDelivery(ctx, d) + } + } +} + +func (w *LikeWorker) 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("like worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) + _ = d.Ack(false) + return + } + log.Printf("like worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) + _ = d.Nack(false, true) + return + } + _ = d.Ack(false) +} + +func (w *LikeWorker) process(ctx context.Context, body []byte) error { + var evt rabbitmq.LikeEvent + if err := json.Unmarshal(body, &evt); err != nil { + // 解析事件失败,直接丢弃 + return nil + } + if evt.UserID == 0 || evt.VideoID == 0 { + return nil + } + + switch evt.Action { + case "like": + return w.applyLike(ctx, evt.UserID, evt.VideoID) + case "unlike": + return w.applyUnlike(ctx, evt.UserID, evt.VideoID) + default: + return nil + } +} + +func (w *LikeWorker) applyLike(ctx context.Context, userID, videoID uint) error { + ok, err := w.videos.IsExist(ctx, videoID) + if err != nil { + return err + } + if !ok { + return nil + } + + created, err := w.likes.LikeIgnoreDuplicate(ctx, &video.Like{ + VideoID: videoID, + AccountID: userID, + CreatedAt: time.Now(), + }) + if err != nil { + return err + } + if !created { + return nil + } + + if err := w.videos.ChangeLikesCount(ctx, videoID, 1); err != nil { + return err + } + return w.videos.ChangePopularity(ctx, videoID, 1) +} + +func (w *LikeWorker) applyUnlike(ctx context.Context, userID, videoID uint) error { + ok, err := w.videos.IsExist(ctx, videoID) + if err != nil { + return err + } + if !ok { + return nil + } + + deleted, err := w.likes.DeleteByVideoAndAccount(ctx, videoID, userID) + if err != nil { + return err + } + if !deleted { + return nil + } + + if err := w.videos.ChangeLikesCount(ctx, videoID, -1); err != nil { + return err + } + return w.videos.ChangePopularity(ctx, videoID, -1) +} diff --git a/backend/internal/worker/popularityworker.go b/backend/internal/worker/popularityworker.go index 8184425..8f1fbaf 100644 --- a/backend/internal/worker/popularityworker.go +++ b/backend/internal/worker/popularityworker.go @@ -1,79 +1,85 @@ -package worker - -import ( - "context" - "encoding/json" - "errors" - "feedsystem_video_go/internal/middleware/rabbitmq" - rediscache "feedsystem_video_go/internal/middleware/redis" - "feedsystem_video_go/internal/video" - "log" - - amqp "github.com/rabbitmq/amqp091-go" -) - -type PopularityWorker struct { - ch *amqp.Channel - cache *rediscache.Client - queue string -} - -func NewPopularityWorker(ch *amqp.Channel, cache *rediscache.Client, queue string) *PopularityWorker { - return &PopularityWorker{ch: ch, cache: cache, queue: queue} -} - -func (w *PopularityWorker) Run(ctx context.Context) error { - if w == nil || w.ch == nil || w.cache == nil { - return errors.New("popularity worker is not initialized") - } - if w.queue == "" { - return errors.New("queue is required") - } - - deliveries, err := w.ch.Consume( - w.queue, - "", - false, - false, - false, - false, - nil, - ) - if err != nil { - return err - } - - for { - select { - case <-ctx.Done(): - return ctx.Err() - case d, ok := <-deliveries: - if !ok { - return errors.New("deliveries channel closed") - } - w.handleDelivery(ctx, d) - } - } -} - -func (w *PopularityWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { - if err := w.process(ctx, d.Body); err != nil { - log.Printf("popularity worker: failed to process message: %v", err) - _ = d.Nack(false, true) - return - } - _ = d.Ack(false) -} - -func (w *PopularityWorker) process(ctx context.Context, body []byte) error { - var evt rabbitmq.PopularityEvent - if err := json.Unmarshal(body, &evt); err != nil { - return nil - } - if evt.VideoID == 0 || evt.Change == 0 { - return nil - } - video.UpdatePopularityCache(ctx, w.cache, evt.VideoID, evt.Change) - return nil -} - +package worker + +import ( + "context" + "encoding/json" + "errors" + "feedsystem_video_go/internal/middleware/rabbitmq" + rediscache "feedsystem_video_go/internal/middleware/redis" + "feedsystem_video_go/internal/video" + "log" + + amqp "github.com/rabbitmq/amqp091-go" +) + +type PopularityWorker struct { + ch *amqp.Channel + cache *rediscache.Client + queue string +} + +func NewPopularityWorker(ch *amqp.Channel, cache *rediscache.Client, queue string) *PopularityWorker { + return &PopularityWorker{ch: ch, cache: cache, queue: queue} +} + +func (w *PopularityWorker) Run(ctx context.Context) error { + if w == nil || w.ch == nil || w.cache == nil { + return errors.New("popularity worker is not initialized") + } + if w.queue == "" { + return errors.New("queue is required") + } + + deliveries, err := w.ch.Consume( + w.queue, + "", + false, + false, + false, + false, + nil, + ) + if err != nil { + return err + } + + for { + select { + case <-ctx.Done(): + return ctx.Err() + case d, ok := <-deliveries: + if !ok { + return errors.New("deliveries channel closed") + } + w.handleDelivery(ctx, d) + } + } +} + +func (w *PopularityWorker) 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("popularity worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) + _ = d.Ack(false) + return + } + log.Printf("popularity worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) + _ = d.Nack(false, true) + return + } + _ = d.Ack(false) +} + +func (w *PopularityWorker) process(ctx context.Context, body []byte) error { + var evt rabbitmq.PopularityEvent + if err := json.Unmarshal(body, &evt); err != nil { + return nil + } + if evt.VideoID == 0 || evt.Change == 0 { + return nil + } + video.UpdatePopularityCache(ctx, w.cache, evt.VideoID, evt.Change) + return nil +} + diff --git a/backend/internal/worker/socialworker.go b/backend/internal/worker/socialworker.go index 58c975a..2aff6ac 100644 --- a/backend/internal/worker/socialworker.go +++ b/backend/internal/worker/socialworker.go @@ -1,101 +1,106 @@ -package worker - -import ( - "context" - "encoding/json" - "errors" - "feedsystem_video_go/internal/middleware/rabbitmq" - "feedsystem_video_go/internal/social" - "log" - - "github.com/go-sql-driver/mysql" - amqp "github.com/rabbitmq/amqp091-go" -) - -type SocialWorker struct { - ch *amqp.Channel - repo *social.SocialRepository - queue string -} - -func NewSocialWorker(ch *amqp.Channel, repo *social.SocialRepository, queue string) *SocialWorker { - return &SocialWorker{ch: ch, repo: repo, queue: queue} -} - -func (w *SocialWorker) Run(ctx context.Context) error { - if w == nil || w.ch == nil || w.repo == nil { - return errors.New("social worker is not initialized") - } - if w.queue == "" { - return errors.New("queue is required") - } - - deliveries, err := w.ch.Consume( - w.queue, - "", - false, - false, - false, - false, - nil, - ) - if err != nil { - return err - } - - for { - select { - case <-ctx.Done(): - return ctx.Err() - case d, ok := <-deliveries: - if !ok { - return errors.New("deliveries channel closed") - } - w.handleDelivery(ctx, d) - } - } -} - -func (w *SocialWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { - if err := w.process(ctx, d.Body); err != nil { - log.Printf("social worker: failed to process message: %v", err) - // 重新入队,稍后重试 - _ = d.Nack(false, true) - return - } - _ = d.Ack(false) -} - -func (w *SocialWorker) process(ctx context.Context, body []byte) error { - var evt rabbitmq.SocialEvent - if err := json.Unmarshal(body, &evt); err != nil { - // 解析事件失败,直接丢弃 - return nil - } - if evt.FollowerID == 0 || evt.VloggerID == 0 { - return nil - } - - switch evt.Action { - case "follow": - err := w.repo.Follow(ctx, &social.Social{ - FollowerID: evt.FollowerID, - VloggerID: evt.VloggerID, - }) - if err == nil { - return nil - } - var mysqlErr *mysql.MySQLError - if errors.As(err, &mysqlErr) && mysqlErr.Number == 1062 { - return nil - } - return err - case "unfollow": - return w.repo.Unfollow(ctx, &social.Social{ - FollowerID: evt.FollowerID, - VloggerID: evt.VloggerID, - }) - default: - return nil - } -} +package worker + +import ( + "context" + "encoding/json" + "errors" + "feedsystem_video_go/internal/middleware/rabbitmq" + "feedsystem_video_go/internal/social" + "log" + + "github.com/go-sql-driver/mysql" + amqp "github.com/rabbitmq/amqp091-go" +) + +type SocialWorker struct { + ch *amqp.Channel + repo *social.SocialRepository + queue string +} + +func NewSocialWorker(ch *amqp.Channel, repo *social.SocialRepository, queue string) *SocialWorker { + return &SocialWorker{ch: ch, repo: repo, queue: queue} +} + +func (w *SocialWorker) Run(ctx context.Context) error { + if w == nil || w.ch == nil || w.repo == nil { + return errors.New("social worker is not initialized") + } + if w.queue == "" { + return errors.New("queue is required") + } + + deliveries, err := w.ch.Consume( + w.queue, + "", + false, + false, + false, + false, + nil, + ) + if err != nil { + return err + } + + for { + select { + case <-ctx.Done(): + return ctx.Err() + case d, ok := <-deliveries: + if !ok { + return errors.New("deliveries channel closed") + } + w.handleDelivery(ctx, d) + } + } +} + +func (w *SocialWorker) 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("social worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) + _ = d.Ack(false) + return + } + log.Printf("social worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) + _ = d.Nack(false, true) + return + } + _ = d.Ack(false) +} + +func (w *SocialWorker) process(ctx context.Context, body []byte) error { + var evt rabbitmq.SocialEvent + if err := json.Unmarshal(body, &evt); err != nil { + // 解析事件失败,直接丢弃 + return nil + } + if evt.FollowerID == 0 || evt.VloggerID == 0 { + return nil + } + + switch evt.Action { + case "follow": + err := w.repo.Follow(ctx, &social.Social{ + FollowerID: evt.FollowerID, + VloggerID: evt.VloggerID, + }) + if err == nil { + return nil + } + var mysqlErr *mysql.MySQLError + if errors.As(err, &mysqlErr) && mysqlErr.Number == 1062 { + return nil + } + return err + case "unfollow": + return w.repo.Unfollow(ctx, &social.Social{ + FollowerID: evt.FollowerID, + VloggerID: evt.VloggerID, + }) + default: + return nil + } +}