feat(P1): MQ Worker 死信队列 — 重试上限3次后 Ack 移入 DLX

This commit is contained in:
Sisyphus
2026-04-25 15:54:55 +08:00
parent d68a4f3f65
commit 9d903ad8e1
6 changed files with 637 additions and 554 deletions

View File

@@ -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)
}

View File

@@ -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
}

View File

@@ -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)
}

View File

@@ -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)
}

View File

@@ -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
}

View File

@@ -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
}
}