| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251 |
- using Application.Abstractions.Data;
- using Application.Abstractions.YouTube;
- using Domain.Entities.Members;
- using Domain.Entities.Members.ValueObject;
- using Microsoft.EntityFrameworkCore;
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Hosting;
- using Microsoft.Extensions.Logging;
- namespace Infrastructure.YouTube;
- /// <summary>
- /// 매일 자정 00:00(KST) 1회 실행 — YouTube 방송 지표 집계.
- /// - 각 활성 채널의 uploads 플레이리스트에서 최근 50개 videoId 조회
- /// - IsFinalized=false 세션 또는 최근 7일 내 업로드 videoId만 대상
- /// - videos.list 로 liveStreamingDetails + statistics 배치 조회
- /// - BroadcastSession 생성/갱신, End+24h 경과 시 IsFinalized=true
- /// - ChannelBroadcastStats 누적 재계산
- ///
- /// API quota: 채널당 1(playlistItems) + 3(videos.list) = 4 units/일
- /// 1,000채널까지 안전 (기본 10,000 quota).
- /// </summary>
- internal sealed class YouTubeDailyAggregatorService(
- IServiceScopeFactory scopeFactory,
- IYouTubeApiService youTubeApi,
- ILogger<YouTubeDailyAggregatorService> logger
- ) : BackgroundService
- {
- private const int BatchSize = 50; // YouTube API 최대 50개/요청
- private const int UploadsLookback = 50; // 채널당 조회할 최근 업로드 수
- private static readonly TimeSpan InitialDelay = TimeSpan.FromMinutes(2);
- // 한국 표준시(KST) 자정 기준으로 실행
- private static readonly TimeZoneInfo Kst = TimeZoneInfo.FindSystemTimeZoneById("Korea Standard Time");
- protected override async Task ExecuteAsync(CancellationToken stoppingToken)
- {
- await Task.Delay(InitialDelay, stoppingToken);
- logger.LogInformation("[YouTubeDailyAggregator] 서비스 시작 — 매일 KST 00:00 실행");
- while (!stoppingToken.IsCancellationRequested)
- {
- var delay = CalculateDelayUntilNextMidnight();
- logger.LogInformation("[YouTubeDailyAggregator] 다음 실행까지 대기: {Hours}시간 {Minutes}분", delay.Hours, delay.Minutes);
- try
- {
- await Task.Delay(delay, stoppingToken);
- }
- catch (TaskCanceledException)
- {
- break;
- }
- try
- {
- await AggregateAllAsync(stoppingToken);
- }
- catch (Exception ex)
- {
- logger.LogError(ex, "[YouTubeDailyAggregator] 집계 실행 중 오류");
- }
- }
- }
- private static TimeSpan CalculateDelayUntilNextMidnight()
- {
- var nowKst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, Kst);
- var nextMidnightKst = nowKst.Date.AddDays(1);
- var delay = nextMidnightKst - nowKst;
- if (delay <= TimeSpan.Zero)
- {
- delay = TimeSpan.FromMinutes(1);
- }
- return delay;
- }
- private async Task AggregateAllAsync(CancellationToken ct)
- {
- logger.LogInformation("[YouTubeDailyAggregator] 집계 시작");
- using var scope = scopeFactory.CreateScope();
- var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
- // YouTubeChannelID 가 설정된 활성 채널만 대상
- var channels = await db.Channel.AsNoTracking()
- .Where(c => c.IsActive && c.YouTubeChannelID != null && c.YouTubeChannelID != "")
- .Select(c => new { c.ID, c.YouTubeChannelID })
- .ToListAsync(ct);
- if (channels.Count == 0)
- {
- logger.LogInformation("[YouTubeDailyAggregator] 대상 채널 없음");
- return;
- }
- var processed = 0;
- var failed = 0;
- foreach (var channel in channels)
- {
- if (ct.IsCancellationRequested)
- {
- break;
- }
- try
- {
- await AggregateChannelAsync(channel.ID, channel.YouTubeChannelID!, ct);
- processed++;
- }
- catch (Exception ex)
- {
- failed++;
- logger.LogError(ex, "[YouTubeDailyAggregator] 채널 {ChannelID} 집계 실패", channel.ID);
- }
- // API rate limit 방어
- await Task.Delay(TimeSpan.FromMilliseconds(500), ct);
- }
- logger.LogInformation("[YouTubeDailyAggregator] 집계 완료 — 성공: {Processed}, 실패: {Failed}", processed, failed);
- }
- private async Task AggregateChannelAsync(int channelID, string youtubeChannelID, CancellationToken ct)
- {
- // 1. 최근 업로드 videoId 조회 (1 unit)
- var recentVideoIds = await youTubeApi.GetRecentUploadsAsync(youtubeChannelID, UploadsLookback, ct);
- if (recentVideoIds.Count == 0)
- {
- return;
- }
- using var scope = scopeFactory.CreateScope();
- var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
- // 2. 기존 세션 중 IsFinalized=false 또는 최근 7일 내 세션 videoId 추가
- var now = DateTime.UtcNow;
- var lookbackStart = now.AddDays(-7);
- var pendingSessions = await db.BroadcastSession
- .Where(s => s.ChannelID == channelID && (!s.IsFinalized || (s.EndAt != null && s.EndAt > lookbackStart)))
- .ToListAsync(ct);
- var targetVideoIds = new HashSet<string>(recentVideoIds);
- foreach (var session in pendingSessions)
- {
- targetVideoIds.Add(session.VideoID);
- }
- // 3. videos.list 배치 조회 (50개씩, 3 units/배치)
- var targetList = targetVideoIds.ToList();
- var allStats = new List<YouTubeVideoStats>();
- for (var i = 0; i < targetList.Count; i += BatchSize)
- {
- var batch = targetList.Skip(i).Take(BatchSize).ToList();
- var stats = await youTubeApi.GetVideoStatsBatchAsync(batch, ct);
- allStats.AddRange(stats);
- if (i + BatchSize < targetList.Count)
- {
- await Task.Delay(TimeSpan.FromMilliseconds(250), ct);
- }
- }
- // 4. 세션 upsert
- var existingByVideoId = pendingSessions.ToDictionary(s => s.VideoID, s => s);
- var sessionsToAdd = new List<BroadcastSession>();
- foreach (var stat in allStats)
- {
- // 라이브 방송이 아닌 일반 업로드는 제외
- if (!stat.IsLiveStream || stat.ActualStartTime == null)
- {
- continue;
- }
- if (existingByVideoId.TryGetValue(stat.VideoID, out var existing))
- {
- existing.UpdateStats(stat.Title, stat.ActualStartTime, stat.ActualEndTime, stat.ViewCount, stat.LikeCount);
- // End + 24h 경과 시 finalize
- if (stat.ActualEndTime.HasValue && stat.ActualEndTime.Value.AddHours(24) < now)
- {
- existing.MarkFinalized();
- }
- }
- else
- {
- var session = BroadcastSession.Create(channelID, BroadcastPlatform.YouTube, stat.VideoID);
- session.UpdateStats(stat.Title, stat.ActualStartTime, stat.ActualEndTime, stat.ViewCount, stat.LikeCount);
- if (stat.ActualEndTime.HasValue && stat.ActualEndTime.Value.AddHours(24) < now)
- {
- session.MarkFinalized();
- }
- sessionsToAdd.Add(session);
- }
- }
- if (sessionsToAdd.Count > 0)
- {
- await db.BroadcastSession.AddRangeAsync(sessionsToAdd, ct);
- }
- await db.SaveChangesAsync(ct);
- // 5. ChannelBroadcastStats 누적 재계산 (모든 세션 스캔)
- await RecalculateChannelStatsAsync(db, channelID, ct);
- }
- private static async Task RecalculateChannelStatsAsync(IAppDbContext db, int channelID, CancellationToken ct)
- {
- var aggregate = await db.BroadcastSession
- .Where(s => s.ChannelID == channelID && s.StartAt != null && s.EndAt != null)
- .GroupBy(s => 1)
- .Select(g => new
- {
- TotalDurationSec = g.Sum(s => (long)s.DurationSec),
- CumulativeViews = g.Sum(s => s.FinalViews),
- TotalLikes = g.Sum(s => s.FinalLikes),
- SessionCount = g.Count(),
- LastBroadcastAt = g.Max(s => s.EndAt)
- })
- .FirstOrDefaultAsync(ct);
- var stats = await db.ChannelBroadcastStats.FirstOrDefaultAsync(s => s.ChannelID == channelID, ct);
- if (aggregate == null)
- {
- return;
- }
- if (stats == null)
- {
- stats = ChannelBroadcastStats.Create(channelID);
- stats.Update(aggregate.TotalDurationSec, aggregate.CumulativeViews, aggregate.TotalLikes, aggregate.SessionCount, aggregate.LastBroadcastAt);
- await db.ChannelBroadcastStats.AddAsync(stats, ct);
- }
- else
- {
- stats.Update(aggregate.TotalDurationSec, aggregate.CumulativeViews, aggregate.TotalLikes, aggregate.SessionCount, aggregate.LastBroadcastAt);
- }
- await db.SaveChangesAsync(ct);
- }
- }
|