feat:时间线MQ发送端和消费端
This commit is contained in:
@@ -24,7 +24,7 @@ func NewDB(dbcfg config.DatabaseConfig) (*gorm.DB, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func AutoMigrate(db *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 {
|
func CloseDB(db *gorm.DB) error {
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
rediscache "feedsystem_video_go/internal/middleware/redis"
|
rediscache "feedsystem_video_go/internal/middleware/redis"
|
||||||
"feedsystem_video_go/internal/social"
|
"feedsystem_video_go/internal/social"
|
||||||
"feedsystem_video_go/internal/video"
|
"feedsystem_video_go/internal/video"
|
||||||
|
"feedsystem_video_go/internal/worker"
|
||||||
"log"
|
"log"
|
||||||
|
|
||||||
"github.com/gin-gonic/gin"
|
"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)
|
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
|
return r
|
||||||
}
|
}
|
||||||
|
|||||||
92
backend/internal/worker/outboxworker.go
Normal file
92
backend/internal/worker/outboxworker.go
Normal file
@@ -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()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user