diff --git a/backend/internal/db/db.go b/backend/internal/db/db.go index f082f97..23593eb 100644 --- a/backend/internal/db/db.go +++ b/backend/internal/db/db.go @@ -24,7 +24,7 @@ func NewDB(dbcfg config.DatabaseConfig) (*gorm.DB, error) { } func AutoMigrate(db *gorm.DB) error { - return db.AutoMigrate(&account.Account{}, &video.Video{}, &video.Like{}, &video.Comment{}, &social.Social{}) + return db.AutoMigrate(&account.Account{}, &video.Video{}, &video.Like{}, &video.Comment{}, &social.Social{}, &video.OutboxMsg{}) } func CloseDB(db *gorm.DB) error { diff --git a/backend/internal/http/router.go b/backend/internal/http/router.go index 8581eb6..7f65f84 100644 --- a/backend/internal/http/router.go +++ b/backend/internal/http/router.go @@ -8,6 +8,7 @@ import ( rediscache "feedsystem_video_go/internal/middleware/redis" "feedsystem_video_go/internal/social" "feedsystem_video_go/internal/video" + "feedsystem_video_go/internal/worker" "log" "github.com/gin-gonic/gin" @@ -127,5 +128,9 @@ func SetRouter(db *gorm.DB, cache *rediscache.Client, rmq *rabbitmq.RabbitMQ) *g { protectedFeedGroup.POST("/listByFollowing", feedHandler.ListByFollowing) } + //worker + timelineMQ, err := rabbitmq.NewTimelineMQ(rmq) + worker.StartOutboxPoller(db, timelineMQ) + worker.StartConsumer(timelineMQ, "video.timeline.update.queue", cache) return r } diff --git a/backend/internal/worker/outboxworker.go b/backend/internal/worker/outboxworker.go new file mode 100644 index 0000000..662b275 --- /dev/null +++ b/backend/internal/worker/outboxworker.go @@ -0,0 +1,92 @@ +package worker + +import ( + "context" + "encoding/json" + "feedsystem_video_go/internal/middleware/rabbitmq" + "feedsystem_video_go/internal/middleware/redis" + "feedsystem_video_go/internal/video" + "fmt" + "log" + "time" + + oredis "github.com/redis/go-redis/v9" + "gorm.io/gorm" +) + +func StartOutboxPoller(db *gorm.DB, tmq *rabbitmq.TimelineMQ) { + go func() { + for { + var messages []video.OutboxMsg + + err := db.Where("status = ?", "pending").Order("create_time ASC").Limit(100).Find(&messages).Error + + if err != nil || len(messages) == 0 { + time.Sleep(1 * time.Second) + continue + } + + for _, msg := range messages { + err := tmq.PublishVideo(context.Background(), msg.VideoID, msg.CreateTime) + + if err == nil { + db.Delete(&msg) + } else { + log.Printf("投递MQ失败: VideoID: %d, err: %v", msg.VideoID, err) + } + } + } + }() +} + +func StartConsumer(tmq *rabbitmq.TimelineMQ, queueName string, redisClient *redis.Client) { + msgs, err := tmq.Ch.Consume( + queueName, + "", + false, + false, + false, + false, + nil, + ) + + if err != nil { + log.Printf("注册消费失败") + } + + go func() { + for msg := range msgs { + var event rabbitmq.TimelineEvent + err := json.Unmarshal(msg.Body, &event) + + if err != nil { + log.Printf("反序列化失败") + msg.Ack(false) + continue + } + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + timelineKey := "feed:global_timeline" + err = redisClient.ZAdd(ctx, timelineKey, oredis.Z{ + Score: float64(event.CreateTime), + Member: fmt.Sprintf("%d", event.ViedoID), + }) + + if err != nil { + log.Printf("写入Zset失败") + msg.Nack(false, true) + cancel() + continue + } + + err = redisClient.ZRemRangeByRank(ctx, timelineKey, 0, -1001) + + if err != nil { + log.Printf("ZRem失败") + } + + msg.Ack(false) + cancel() + } + }() +}