Files
VLoop/backend/internal/feed/service.go
2025-12-25 22:40:20 +08:00

257 lines
7.7 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package feed
import (
"context"
"encoding/json"
rediscache "feedsystem_video_go/internal/redis"
"feedsystem_video_go/internal/video"
"fmt"
"time"
)
type FeedService struct {
repo *FeedRepository
likeRepo *video.LikeRepository
cache *rediscache.Client
cacheTTL time.Duration
}
func NewFeedService(repo *FeedRepository, likeRepo *video.LikeRepository, cache *rediscache.Client) *FeedService {
return &FeedService{repo: repo, likeRepo: likeRepo, cache: cache, cacheTTL: 5 * time.Second}
}
// 查询最新视频
func (f *FeedService) ListLatest(ctx context.Context, limit int, latestBefore time.Time, viewerAccountID uint) (ListLatestResponse, error) {
// 从数据库中查询最新视频
doListLatestFromDB := func() (ListLatestResponse, error) {
videos, err := f.repo.ListLatest(ctx, limit, latestBefore)
if err != nil {
return ListLatestResponse{}, err
}
var nextTime int64
if len(videos) > 0 {
nextTime = videos[len(videos)-1].CreateTime.Unix()
} else {
nextTime = 0
}
hasMore := len(videos) == limit
feedVideos, err := f.buildFeedVideos(ctx, videos, viewerAccountID)
if err != nil {
return ListLatestResponse{}, err
}
resp := ListLatestResponse{
VideoList: feedVideos,
NextTime: nextTime,
HasMore: hasMore,
}
return resp, nil
}
// 先从缓存中查询
var cacheKey string
if viewerAccountID == 0 && f.cache != nil {
before := int64(0)
if !latestBefore.IsZero() {
before = latestBefore.Unix()
}
cacheKey = fmt.Sprintf("feed:listLatest:limit=%d:before=%d", limit, before)
cacheCtx, cancel := context.WithTimeout(ctx, 50*time.Millisecond)
defer cancel()
b, err := f.cache.GetBytes(cacheCtx, cacheKey)
if err == nil {
var cached ListLatestResponse
if err := json.Unmarshal(b, &cached); err == nil {
return cached, nil
}
} else if rediscache.IsMiss(err) { // 缓存未命中
lockKey := "lock:" + cacheKey
// 缓存未命中,尝试加锁
token, locked, _ := f.cache.Lock(cacheCtx, lockKey, 500*time.Millisecond)
if locked {
defer func() { _ = f.cache.Unlock(context.Background(), lockKey, token) }()
if b, err := f.cache.GetBytes(cacheCtx, cacheKey); err == nil {
var cached ListLatestResponse
if err := json.Unmarshal(b, &cached); err == nil {
return cached, nil
}
} else { // 缓存未命中,从数据库中查询
resp, err := doListLatestFromDB()
if err != nil {
return ListLatestResponse{}, err
}
if b, err := json.Marshal(resp); err == nil {
_ = f.cache.SetBytes(cacheCtx, cacheKey, b, f.cacheTTL)
}
return resp, nil
}
} else { // 缓存未命中其他goroutine正在查询等待
for i := 0; i < 5; i++ {
time.Sleep(20 * time.Millisecond)
if b, err := f.cache.GetBytes(cacheCtx, cacheKey); err == nil {
var cached ListLatestResponse
if err := json.Unmarshal(b, &cached); err == nil {
return cached, nil
}
}
}
}
}
}
// 缓存中没有查询到结果,从数据库中查询
resp, err := doListLatestFromDB()
if err != nil {
return ListLatestResponse{}, err
}
// 缓存查询结果
if cacheKey != "" {
if b, err := json.Marshal(resp); err == nil {
cacheCtx, cancel := context.WithTimeout(ctx, 50*time.Millisecond)
defer cancel()
_ = f.cache.SetBytes(cacheCtx, cacheKey, b, f.cacheTTL)
}
}
return resp, nil
}
// 按照点赞数查询视频
func (f *FeedService) ListLikesCount(ctx context.Context, limit int, cursor *LikesCountCursor, viewerAccountID uint) (ListLikesCountResponse, error) {
videos, err := f.repo.ListLikesCountWithCursor(ctx, limit, cursor)
if err != nil {
return ListLikesCountResponse{}, err
}
hasMore := len(videos) == limit
feedVideos, err := f.buildFeedVideos(ctx, videos, viewerAccountID)
if err != nil {
return ListLikesCountResponse{}, err
}
resp := ListLikesCountResponse{
VideoList: feedVideos,
HasMore: hasMore,
}
if len(videos) > 0 {
last := videos[len(videos)-1]
nextLikesCountBefore := last.LikesCount
nextIDBefore := last.ID
resp.NextLikesCountBefore = &nextLikesCountBefore
resp.NextIDBefore = &nextIDBefore
}
return resp, nil
}
// 按照关注列表查询视频
func (f *FeedService) ListByFollowing(ctx context.Context, limit int, latestBefore time.Time, viewerAccountID uint) (ListByFollowingResponse, error) {
doListByFollowingFromDB := func() (ListByFollowingResponse, error) {
videos, err := f.repo.ListByFollowing(ctx, limit, viewerAccountID, latestBefore)
if err != nil {
return ListByFollowingResponse{}, err
}
var nextTime int64
if len(videos) > 0 {
nextTime = videos[len(videos)-1].CreateTime.Unix()
} else {
nextTime = 0
}
hasMore := len(videos) == limit
feedVideos, err := f.buildFeedVideos(ctx, videos, viewerAccountID)
if err != nil {
return ListByFollowingResponse{}, err
}
resp := ListByFollowingResponse{
VideoList: feedVideos,
NextTime: nextTime,
HasMore: hasMore,
}
return resp, nil
}
var cacheKey string
if viewerAccountID != 0 && f.cache != nil {
before := int64(0)
if !latestBefore.IsZero() {
before = latestBefore.Unix()
}
cacheKey = fmt.Sprintf("feed:listByFollowing:limit=%d:accountID=%d:before=%d", limit, viewerAccountID, before)
cacheCtx, cancel := context.WithTimeout(ctx, 50*time.Millisecond)
defer cancel()
b, err := f.cache.GetBytes(cacheCtx, cacheKey)
if err == nil {
var cached ListByFollowingResponse
if err := json.Unmarshal(b, &cached); err == nil {
return cached, nil
}
} else if rediscache.IsMiss(err) { // 缓存未命中
lockKey := "lock:" + cacheKey
// 缓存未命中,尝试加锁
token, locked, _ := f.cache.Lock(cacheCtx, lockKey, 500*time.Millisecond)
if locked {
defer func() { _ = f.cache.Unlock(context.Background(), lockKey, token) }()
if b, err := f.cache.GetBytes(cacheCtx, cacheKey); err == nil {
var cached ListByFollowingResponse
if err := json.Unmarshal(b, &cached); err == nil {
return cached, nil
}
} else { // 缓存未命中,从数据库中查询
resp, err := doListByFollowingFromDB()
if err != nil {
return ListByFollowingResponse{}, err
}
if b, err := json.Marshal(resp); err == nil {
_ = f.cache.SetBytes(cacheCtx, cacheKey, b, f.cacheTTL)
}
return resp, nil
}
} else {
for i := 0; i < 5; i++ {
time.Sleep(20 * time.Millisecond)
if b, err := f.cache.GetBytes(cacheCtx, cacheKey); err == nil {
var cached ListByFollowingResponse
if err := json.Unmarshal(b, &cached); err == nil {
return cached, nil
}
}
}
}
}
}
resp, err := doListByFollowingFromDB()
if err != nil {
return ListByFollowingResponse{}, err
}
if cacheKey != "" {
if b, err := json.Marshal(resp); err == nil {
cacheCtx, cancel := context.WithTimeout(ctx, 50*time.Millisecond)
defer cancel()
_ = f.cache.SetBytes(cacheCtx, cacheKey, b, f.cacheTTL)
}
}
return resp, nil
}
func (f *FeedService) buildFeedVideos(ctx context.Context, videos []*video.Video, viewerAccountID uint) ([]FeedVideoItem, error) {
feedVideos := make([]FeedVideoItem, 0, len(videos))
videoIDs := make([]uint, len(videos))
for i, v := range videos {
videoIDs[i] = v.ID
}
likedMap, err := f.likeRepo.BatchGetLiked(ctx, videoIDs, viewerAccountID)
if err != nil {
return nil, err
}
for _, video := range videos {
feedVideos = append(feedVideos, FeedVideoItem{
ID: video.ID,
Author: FeedAuthor{ID: video.AuthorID, Username: video.Username},
Title: video.Title,
Description: video.Description,
PlayURL: video.PlayURL,
CoverURL: video.CoverURL,
CreateTime: video.CreateTime.Unix(),
LikesCount: video.LikesCount,
IsLiked: likedMap[video.ID],
})
}
return feedVideos, nil
}