fix: Worker 进程独立 Channel + 进程内指数退避重试

- 每个 Worker 消费者使用独立 AMQP Channel,互不影响
- 用 runWorkerWithRetry 包装每个 Worker,断开自动重连
- 用进程内指数退避重试替代失效的 Nack+DLX 重试机制
- 拓扑声明改用临时 Channel,声明完即关闭
This commit is contained in:
yiyiis
2026-05-23 09:23:14 +08:00
parent 33c0bd418c
commit 589116b141
6 changed files with 154 additions and 74 deletions

View File

@@ -55,6 +55,38 @@ func connectWithRetry(name string, maxRetries int, fn func() error) {
log.Fatalf("%s: 超过最大重试次数", name) 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() { func main() {
// 加载 .env本地开发 // 加载 .env本地开发
if err := godotenv.Load(); err != nil { if err := godotenv.Load(); err != nil {
@@ -110,42 +142,33 @@ func main() {
return err return err
}) })
defer conn.Close() defer conn.Close()
// 创建 RabbitMQ 通道
ch, err := conn.Channel() // 用临时 Channel 声明拓扑(持久化队列,声明一次即可)
topoCh, err := conn.Channel()
if err != nil { if err != nil {
log.Fatalf("Failed to open rabbitmq channel: %v", err) log.Fatalf("Failed to open topology channel: %v", err)
} }
defer ch.Close() if err := declareSocialTopology(topoCh); err != nil {
// 声明 Social 交换机和队列
if err := declareSocialTopology(ch); err != nil {
log.Fatalf("Failed to declare social topology: %v", err) 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) 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) log.Fatalf("Failed to declare comment topology: %v", err)
} }
if cache != nil { if cache != nil {
if err := declarePopularityTopology(ch); err != nil { if err := declarePopularityTopology(topoCh); err != nil {
log.Fatalf("Failed to declare popularity topology: %v", err) log.Fatalf("Failed to declare popularity topology: %v", err)
} }
} }
if err := ch.Qos(50, 0, false); err != nil { topoCh.Close()
log.Fatalf("Failed to set qos: %v", err)
}
repo := social.NewSocialRepository(sqlDB) // 准备 repo
socialWorker := worker.NewSocialWorker(ch, repo, socialQueue) socialRepo := social.NewSocialRepository(sqlDB)
videoRepo := video.NewVideoRepository(sqlDB) videoRepo := video.NewVideoRepository(sqlDB)
likeRepo := video.NewLikeRepository(sqlDB) likeRepo := video.NewLikeRepository(sqlDB)
commentRepo := video.NewCommentRepository(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) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop() defer stop()
@@ -162,22 +185,26 @@ func main() {
defer pprofServer.Close() defer pprofServer.Close()
} }
errCh := make(chan error, 4) // 每个 Worker 独立 Channel + 自动重连
log.Printf("Worker started, consuming queue=%s", socialQueue) go runWorkerWithRetry(ctx, "SocialWorker", conn, func(ch *amqp.Channel) error {
go func() { errCh <- socialWorker.Run(ctx) }() return worker.NewSocialWorker(ch, socialRepo, socialQueue).Run(ctx)
log.Printf("Worker started, consuming queue=%s", likeQueue) })
go func() { errCh <- likeWorker.Run(ctx) }() go runWorkerWithRetry(ctx, "LikeWorker", conn, func(ch *amqp.Channel) error {
log.Printf("Worker started, consuming queue=%s", commentQueue) return worker.NewLikeWorker(ch, likeRepo, videoRepo, likeQueue).Run(ctx)
go func() { errCh <- commentWorker.Run(ctx) }() })
if popularityWorker != nil { go runWorkerWithRetry(ctx, "CommentWorker", conn, func(ch *amqp.Channel) error {
log.Printf("Worker started, consuming queue=%s", popularityQueue) return worker.NewCommentWorker(ch, commentRepo, videoRepo, commentQueue).Run(ctx)
go func() { errCh <- popularityWorker.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 { <-ctx.Done()
log.Fatalf("Worker stopped: %v", err) log.Printf("Worker shutting down...")
} time.Sleep(2 * time.Second) // 等待正在处理的消息完成
log.Printf("Worker stopped") log.Printf("Worker stopped")
} }

View File

@@ -8,6 +8,7 @@ import (
"feedsystem_video_go/internal/video" "feedsystem_video_go/internal/video"
"log" "log"
"strings" "strings"
"time"
amqp "github.com/rabbitmq/amqp091-go" 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) { func (w *CommentWorker) handleDelivery(ctx context.Context, d amqp.Delivery) {
if err := w.process(ctx, d.Body); err != nil { const maxRetries = 3
retryCount := rabbitmq.GetRetryCount(d) for i := 0; i <= maxRetries; i++ {
if retryCount >= rabbitmq.MaxRetryCount { select {
log.Printf("comment worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) case <-ctx.Done():
_ = d.Ack(false) _ = d.Nack(false, true)
return return
default:
} }
log.Printf("comment worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) if err := w.process(ctx, d.Body); err != nil {
_ = d.Nack(false, true) if i >= maxRetries {
log.Printf("comment worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err)
_ = d.Ack(false)
return
}
wait := time.Duration(1<<uint(i)) * time.Second
log.Printf("comment worker: 处理失败, %v 后重试 (%d/%d): %v", wait, i+1, maxRetries, err)
time.Sleep(wait)
continue
}
_ = d.Ack(false)
return return
} }
_ = d.Ack(false)
} }
func (w *CommentWorker) process(ctx context.Context, body []byte) error { func (w *CommentWorker) process(ctx context.Context, body []byte) error {

View File

@@ -58,18 +58,28 @@ func (w *LikeWorker) Run(ctx context.Context) error {
} }
func (w *LikeWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { func (w *LikeWorker) handleDelivery(ctx context.Context, d amqp.Delivery) {
if err := w.process(ctx, d.Body); err != nil { const maxRetries = 3
retryCount := rabbitmq.GetRetryCount(d) for i := 0; i <= maxRetries; i++ {
if retryCount >= rabbitmq.MaxRetryCount { select {
log.Printf("like worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) case <-ctx.Done():
_ = d.Ack(false) _ = d.Nack(false, true)
return return
default:
} }
log.Printf("like worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) if err := w.process(ctx, d.Body); err != nil {
_ = d.Nack(false, true) if i >= maxRetries {
log.Printf("like worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err)
_ = d.Ack(false)
return
}
wait := time.Duration(1<<uint(i)) * time.Second
log.Printf("like worker: 处理失败, %v 后重试 (%d/%d): %v", wait, i+1, maxRetries, err)
time.Sleep(wait)
continue
}
_ = d.Ack(false)
return return
} }
_ = d.Ack(false)
} }
func (w *LikeWorker) process(ctx context.Context, body []byte) error { func (w *LikeWorker) process(ctx context.Context, body []byte) error {

View File

@@ -66,18 +66,28 @@ func (w *NotificationWorker) Run(ctx context.Context) error {
} }
func (w *NotificationWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { func (w *NotificationWorker) handleDelivery(ctx context.Context, d amqp.Delivery) {
retryCount := rabbitmq.GetRetryCount(d) const maxRetries = 3
if err := w.process(ctx, d); err != nil { for i := 0; i <= maxRetries; i++ {
if retryCount >= rabbitmq.MaxRetryCount { select {
log.Printf("notification worker: max retries, dropping: %v", err) case <-ctx.Done():
_ = d.Ack(false) _ = d.Nack(false, true)
return return
default:
} }
log.Printf("notification worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) if err := w.process(ctx, d); err != nil {
_ = d.Nack(false, true) if i >= maxRetries {
log.Printf("notification worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err)
_ = d.Ack(false)
return
}
wait := time.Duration(1<<uint(i)) * time.Second
log.Printf("notification worker: 处理失败, %v 后重试 (%d/%d): %v", wait, i+1, maxRetries, err)
time.Sleep(wait)
continue
}
_ = d.Ack(false)
return return
} }
_ = d.Ack(false)
} }
func (w *NotificationWorker) process(ctx context.Context, d amqp.Delivery) error { func (w *NotificationWorker) process(ctx context.Context, d amqp.Delivery) error {

View File

@@ -8,6 +8,7 @@ import (
rediscache "feedsystem_video_go/internal/middleware/redis" rediscache "feedsystem_video_go/internal/middleware/redis"
"feedsystem_video_go/internal/video" "feedsystem_video_go/internal/video"
"log" "log"
"time"
amqp "github.com/rabbitmq/amqp091-go" amqp "github.com/rabbitmq/amqp091-go"
) )
@@ -57,18 +58,28 @@ func (w *PopularityWorker) Run(ctx context.Context) error {
} }
func (w *PopularityWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { func (w *PopularityWorker) handleDelivery(ctx context.Context, d amqp.Delivery) {
if err := w.process(ctx, d.Body); err != nil { const maxRetries = 3
retryCount := rabbitmq.GetRetryCount(d) for i := 0; i <= maxRetries; i++ {
if retryCount >= rabbitmq.MaxRetryCount { select {
log.Printf("popularity worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) case <-ctx.Done():
_ = d.Ack(false) _ = d.Nack(false, true)
return return
default:
} }
log.Printf("popularity worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) if err := w.process(ctx, d.Body); err != nil {
_ = d.Nack(false, true) if i >= maxRetries {
log.Printf("popularity worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err)
_ = d.Ack(false)
return
}
wait := time.Duration(1<<uint(i)) * time.Second
log.Printf("popularity worker: 处理失败, %v 后重试 (%d/%d): %v", wait, i+1, maxRetries, err)
time.Sleep(wait)
continue
}
_ = d.Ack(false)
return return
} }
_ = d.Ack(false)
} }
func (w *PopularityWorker) process(ctx context.Context, body []byte) error { func (w *PopularityWorker) process(ctx context.Context, body []byte) error {

View File

@@ -7,6 +7,7 @@ import (
"feedsystem_video_go/internal/middleware/rabbitmq" "feedsystem_video_go/internal/middleware/rabbitmq"
"feedsystem_video_go/internal/social" "feedsystem_video_go/internal/social"
"log" "log"
"time"
"github.com/go-sql-driver/mysql" "github.com/go-sql-driver/mysql"
amqp "github.com/rabbitmq/amqp091-go" amqp "github.com/rabbitmq/amqp091-go"
@@ -57,18 +58,28 @@ func (w *SocialWorker) Run(ctx context.Context) error {
} }
func (w *SocialWorker) handleDelivery(ctx context.Context, d amqp.Delivery) { func (w *SocialWorker) handleDelivery(ctx context.Context, d amqp.Delivery) {
if err := w.process(ctx, d.Body); err != nil { const maxRetries = 3
retryCount := rabbitmq.GetRetryCount(d) for i := 0; i <= maxRetries; i++ {
if retryCount >= rabbitmq.MaxRetryCount { select {
log.Printf("social worker: max retries exceeded (%d), moving to DLX: %v", retryCount, err) case <-ctx.Done():
_ = d.Ack(false) _ = d.Nack(false, true)
return return
default:
} }
log.Printf("social worker: failed (retry %d/%d): %v", retryCount+1, rabbitmq.MaxRetryCount, err) if err := w.process(ctx, d.Body); err != nil {
_ = d.Nack(false, true) if i >= maxRetries {
log.Printf("social worker: 重试 %d 次后仍失败, 丢弃: %v", maxRetries, err)
_ = d.Ack(false)
return
}
wait := time.Duration(1<<uint(i)) * time.Second
log.Printf("social worker: 处理失败, %v 后重试 (%d/%d): %v", wait, i+1, maxRetries, err)
time.Sleep(wait)
continue
}
_ = d.Ack(false)
return return
} }
_ = d.Ack(false)
} }
func (w *SocialWorker) process(ctx context.Context, body []byte) error { func (w *SocialWorker) process(ctx context.Context, body []byte) error {