YouTubeDailyAggregatorService.cs 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251
  1. using Application.Abstractions.Data;
  2. using Application.Abstractions.YouTube;
  3. using Domain.Entities.Members;
  4. using Domain.Entities.Members.ValueObject;
  5. using Microsoft.EntityFrameworkCore;
  6. using Microsoft.Extensions.DependencyInjection;
  7. using Microsoft.Extensions.Hosting;
  8. using Microsoft.Extensions.Logging;
  9. namespace Infrastructure.YouTube;
  10. /// <summary>
  11. /// 매일 자정 00:00(KST) 1회 실행 — YouTube 방송 지표 집계.
  12. /// - 각 활성 채널의 uploads 플레이리스트에서 최근 50개 videoId 조회
  13. /// - IsFinalized=false 세션 또는 최근 7일 내 업로드 videoId만 대상
  14. /// - videos.list 로 liveStreamingDetails + statistics 배치 조회
  15. /// - BroadcastSession 생성/갱신, End+24h 경과 시 IsFinalized=true
  16. /// - ChannelBroadcastStats 누적 재계산
  17. ///
  18. /// API quota: 채널당 1(playlistItems) + 3(videos.list) = 4 units/일
  19. /// 1,000채널까지 안전 (기본 10,000 quota).
  20. /// </summary>
  21. internal sealed class YouTubeDailyAggregatorService(
  22. IServiceScopeFactory scopeFactory,
  23. IYouTubeApiService youTubeApi,
  24. ILogger<YouTubeDailyAggregatorService> logger
  25. ) : BackgroundService
  26. {
  27. private const int BatchSize = 50; // YouTube API 최대 50개/요청
  28. private const int UploadsLookback = 50; // 채널당 조회할 최근 업로드 수
  29. private static readonly TimeSpan InitialDelay = TimeSpan.FromMinutes(2);
  30. // 한국 표준시(KST) 자정 기준으로 실행
  31. private static readonly TimeZoneInfo Kst = TimeZoneInfo.FindSystemTimeZoneById("Korea Standard Time");
  32. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  33. {
  34. await Task.Delay(InitialDelay, stoppingToken);
  35. logger.LogInformation("[YouTubeDailyAggregator] 서비스 시작 — 매일 KST 00:00 실행");
  36. while (!stoppingToken.IsCancellationRequested)
  37. {
  38. var delay = CalculateDelayUntilNextMidnight();
  39. logger.LogInformation("[YouTubeDailyAggregator] 다음 실행까지 대기: {Hours}시간 {Minutes}분", delay.Hours, delay.Minutes);
  40. try
  41. {
  42. await Task.Delay(delay, stoppingToken);
  43. }
  44. catch (TaskCanceledException)
  45. {
  46. break;
  47. }
  48. try
  49. {
  50. await AggregateAllAsync(stoppingToken);
  51. }
  52. catch (Exception ex)
  53. {
  54. logger.LogError(ex, "[YouTubeDailyAggregator] 집계 실행 중 오류");
  55. }
  56. }
  57. }
  58. private static TimeSpan CalculateDelayUntilNextMidnight()
  59. {
  60. var nowKst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, Kst);
  61. var nextMidnightKst = nowKst.Date.AddDays(1);
  62. var delay = nextMidnightKst - nowKst;
  63. if (delay <= TimeSpan.Zero)
  64. {
  65. delay = TimeSpan.FromMinutes(1);
  66. }
  67. return delay;
  68. }
  69. private async Task AggregateAllAsync(CancellationToken ct)
  70. {
  71. logger.LogInformation("[YouTubeDailyAggregator] 집계 시작");
  72. using var scope = scopeFactory.CreateScope();
  73. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  74. // YouTubeChannelID 가 설정된 활성 채널만 대상
  75. var channels = await db.Channel.AsNoTracking()
  76. .Where(c => c.IsActive && c.YouTubeChannelID != null && c.YouTubeChannelID != "")
  77. .Select(c => new { c.ID, c.YouTubeChannelID })
  78. .ToListAsync(ct);
  79. if (channels.Count == 0)
  80. {
  81. logger.LogInformation("[YouTubeDailyAggregator] 대상 채널 없음");
  82. return;
  83. }
  84. var processed = 0;
  85. var failed = 0;
  86. foreach (var channel in channels)
  87. {
  88. if (ct.IsCancellationRequested)
  89. {
  90. break;
  91. }
  92. try
  93. {
  94. await AggregateChannelAsync(channel.ID, channel.YouTubeChannelID!, ct);
  95. processed++;
  96. }
  97. catch (Exception ex)
  98. {
  99. failed++;
  100. logger.LogError(ex, "[YouTubeDailyAggregator] 채널 {ChannelID} 집계 실패", channel.ID);
  101. }
  102. // API rate limit 방어
  103. await Task.Delay(TimeSpan.FromMilliseconds(500), ct);
  104. }
  105. logger.LogInformation("[YouTubeDailyAggregator] 집계 완료 — 성공: {Processed}, 실패: {Failed}", processed, failed);
  106. }
  107. private async Task AggregateChannelAsync(int channelID, string youtubeChannelID, CancellationToken ct)
  108. {
  109. // 1. 최근 업로드 videoId 조회 (1 unit)
  110. var recentVideoIds = await youTubeApi.GetRecentUploadsAsync(youtubeChannelID, UploadsLookback, ct);
  111. if (recentVideoIds.Count == 0)
  112. {
  113. return;
  114. }
  115. using var scope = scopeFactory.CreateScope();
  116. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  117. // 2. 기존 세션 중 IsFinalized=false 또는 최근 7일 내 세션 videoId 추가
  118. var now = DateTime.UtcNow;
  119. var lookbackStart = now.AddDays(-7);
  120. var pendingSessions = await db.BroadcastSession
  121. .Where(s => s.ChannelID == channelID && (!s.IsFinalized || (s.EndAt != null && s.EndAt > lookbackStart)))
  122. .ToListAsync(ct);
  123. var targetVideoIds = new HashSet<string>(recentVideoIds);
  124. foreach (var session in pendingSessions)
  125. {
  126. targetVideoIds.Add(session.VideoID);
  127. }
  128. // 3. videos.list 배치 조회 (50개씩, 3 units/배치)
  129. var targetList = targetVideoIds.ToList();
  130. var allStats = new List<YouTubeVideoStats>();
  131. for (var i = 0; i < targetList.Count; i += BatchSize)
  132. {
  133. var batch = targetList.Skip(i).Take(BatchSize).ToList();
  134. var stats = await youTubeApi.GetVideoStatsBatchAsync(batch, ct);
  135. allStats.AddRange(stats);
  136. if (i + BatchSize < targetList.Count)
  137. {
  138. await Task.Delay(TimeSpan.FromMilliseconds(250), ct);
  139. }
  140. }
  141. // 4. 세션 upsert
  142. var existingByVideoId = pendingSessions.ToDictionary(s => s.VideoID, s => s);
  143. var sessionsToAdd = new List<BroadcastSession>();
  144. foreach (var stat in allStats)
  145. {
  146. // 라이브 방송이 아닌 일반 업로드는 제외
  147. if (!stat.IsLiveStream || stat.ActualStartTime == null)
  148. {
  149. continue;
  150. }
  151. if (existingByVideoId.TryGetValue(stat.VideoID, out var existing))
  152. {
  153. existing.UpdateStats(stat.Title, stat.ActualStartTime, stat.ActualEndTime, stat.ViewCount, stat.LikeCount);
  154. // End + 24h 경과 시 finalize
  155. if (stat.ActualEndTime.HasValue && stat.ActualEndTime.Value.AddHours(24) < now)
  156. {
  157. existing.MarkFinalized();
  158. }
  159. }
  160. else
  161. {
  162. var session = BroadcastSession.Create(channelID, BroadcastPlatform.YouTube, stat.VideoID);
  163. session.UpdateStats(stat.Title, stat.ActualStartTime, stat.ActualEndTime, stat.ViewCount, stat.LikeCount);
  164. if (stat.ActualEndTime.HasValue && stat.ActualEndTime.Value.AddHours(24) < now)
  165. {
  166. session.MarkFinalized();
  167. }
  168. sessionsToAdd.Add(session);
  169. }
  170. }
  171. if (sessionsToAdd.Count > 0)
  172. {
  173. await db.BroadcastSession.AddRangeAsync(sessionsToAdd, ct);
  174. }
  175. await db.SaveChangesAsync(ct);
  176. // 5. ChannelBroadcastStats 누적 재계산 (모든 세션 스캔)
  177. await RecalculateChannelStatsAsync(db, channelID, ct);
  178. }
  179. private static async Task RecalculateChannelStatsAsync(IAppDbContext db, int channelID, CancellationToken ct)
  180. {
  181. var aggregate = await db.BroadcastSession
  182. .Where(s => s.ChannelID == channelID && s.StartAt != null && s.EndAt != null)
  183. .GroupBy(s => 1)
  184. .Select(g => new
  185. {
  186. TotalDurationSec = g.Sum(s => (long)s.DurationSec),
  187. CumulativeViews = g.Sum(s => s.FinalViews),
  188. TotalLikes = g.Sum(s => s.FinalLikes),
  189. SessionCount = g.Count(),
  190. LastBroadcastAt = g.Max(s => s.EndAt)
  191. })
  192. .FirstOrDefaultAsync(ct);
  193. var stats = await db.ChannelBroadcastStats.FirstOrDefaultAsync(s => s.ChannelID == channelID, ct);
  194. if (aggregate == null)
  195. {
  196. return;
  197. }
  198. if (stats == null)
  199. {
  200. stats = ChannelBroadcastStats.Create(channelID);
  201. stats.Update(aggregate.TotalDurationSec, aggregate.CumulativeViews, aggregate.TotalLikes, aggregate.SessionCount, aggregate.LastBroadcastAt);
  202. await db.ChannelBroadcastStats.AddAsync(stats, ct);
  203. }
  204. else
  205. {
  206. stats.Update(aggregate.TotalDurationSec, aggregate.CumulativeViews, aggregate.TotalLikes, aggregate.SessionCount, aggregate.LastBroadcastAt);
  207. }
  208. await db.SaveChangesAsync(ct);
  209. }
  210. }