YouTubeLiveViewerPoller.cs 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171
  1. using Application.Abstractions.Data;
  2. using Application.Abstractions.Hub;
  3. using Application.Abstractions.YouTube;
  4. using Microsoft.EntityFrameworkCore;
  5. using Microsoft.Extensions.DependencyInjection;
  6. using Microsoft.Extensions.Hosting;
  7. using Microsoft.Extensions.Logging;
  8. using StackExchange.Redis;
  9. namespace Infrastructure.YouTube;
  10. /// <summary>
  11. /// 30초 주기로 라이브 중인 채널의 concurrentViewers 를 일괄 조회하여 Redis 에 갱신.
  12. ///
  13. /// Quota 예산:
  14. /// - videos.list (id=v1,v2,...50) = 1 unit / 50 videos
  15. /// - 30초 주기 + batch 50 = 분당 최대 2 call = 하루 최대 2,880 unit (기본 quota 10,000 대비 29%)
  16. /// - 라이브 채널 50개 미만이면 비용 동일 (1 batch 내 처리)
  17. ///
  18. /// 분산 락:
  19. /// - Redis SETNX 기반 락으로 멀티 인스턴스 중복 실행 방지 (quota 2배 소모 방지)
  20. /// - 락 TTL 은 PollInterval 보다 약간 짧게 설정
  21. ///
  22. /// SignalR 연동:
  23. /// - 시청자 수 갱신마다 ReceiveChannelStatus 브로드캐스트 → ChannelSidebar / WatchView 실시간 반영
  24. /// - concurrentViewers 가 응답에 없으면 라이브 종료로 판단 → ClearLiveAsync + IsLive=false 브로드캐스트
  25. /// </summary>
  26. internal sealed class YouTubeLiveViewerPoller(
  27. IYouTubeLiveStateStore liveStateStore,
  28. IYouTubeApiService youTubeApi,
  29. IChannelStatusBroadcaster statusBroadcaster,
  30. IServiceScopeFactory scopeFactory,
  31. IConnectionMultiplexer redis,
  32. ILogger<YouTubeLiveViewerPoller> logger
  33. ) : BackgroundService
  34. {
  35. private const string LockKey = "youtube:viewer-poller:lock";
  36. private static readonly TimeSpan InitialDelay = TimeSpan.FromSeconds(30);
  37. private static readonly TimeSpan PollInterval = TimeSpan.FromSeconds(30);
  38. private static readonly TimeSpan LockDuration = TimeSpan.FromSeconds(25);
  39. private readonly string _instanceId = Guid.NewGuid().ToString("N");
  40. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  41. {
  42. await Task.Delay(InitialDelay, stoppingToken);
  43. logger.LogInformation("[ViewerPoller] 서비스 시작 — 폴링 주기: {Interval}초, 인스턴스: {InstanceId}",
  44. PollInterval.TotalSeconds, _instanceId);
  45. while (!stoppingToken.IsCancellationRequested)
  46. {
  47. try
  48. {
  49. await PollOnceAsync(stoppingToken);
  50. }
  51. catch (OperationCanceledException)
  52. {
  53. break;
  54. }
  55. catch (Exception ex)
  56. {
  57. logger.LogError(ex, "[ViewerPoller] Loop 예외");
  58. }
  59. try
  60. {
  61. await Task.Delay(PollInterval, stoppingToken);
  62. }
  63. catch (OperationCanceledException)
  64. {
  65. break;
  66. }
  67. }
  68. }
  69. private async Task PollOnceAsync(CancellationToken ct)
  70. {
  71. var db = redis.GetDatabase();
  72. // 분산 락 획득 (멀티 인스턴스 환경에서 중복 실행 방지)
  73. var acquired = await db.StringSetAsync(LockKey, _instanceId, LockDuration, When.NotExists);
  74. if (!acquired)
  75. {
  76. logger.LogDebug("[ViewerPoller] 다른 인스턴스가 락 보유 중 — skip");
  77. return;
  78. }
  79. try
  80. {
  81. var allLive = await liveStateStore.GetAllLiveAsync();
  82. var liveList = allLive.Where(l => l.IsLive).ToList();
  83. if (liveList.Count == 0)
  84. {
  85. logger.LogDebug("[ViewerPoller] 라이브 채널 없음 — skip");
  86. return;
  87. }
  88. var videoIds = liveList.Select(l => l.VideoId).ToList();
  89. var viewers = await youTubeApi.GetConcurrentViewersAsync(videoIds, ct);
  90. // 🛡️ Quota 초과 방어: viewers 응답이 통째로 비어있으면 API 실패(quotaExceeded 등)로 판단,
  91. // ended 오인 처리 안 하고 이번 cycle skip. 다음 cycle 또는 다음 PubSub에서 재처리됨.
  92. // 단일 video만 missing이면 진짜 ended로 처리 (정상 분기).
  93. if (viewers.Count == 0)
  94. {
  95. logger.LogWarning("[ViewerPoller] viewers 응답 empty (live={LiveCount}) — API 실패 추정 (quota?). ended 오인 방지 위해 skip",
  96. liveList.Count);
  97. return;
  98. }
  99. // YouTube channelId → 내부 Channel.SID 매핑 (SignalR 브로드캐스트용)
  100. var ytChannelIds = liveList.Select(l => l.ChannelId).Distinct().ToList();
  101. using var scope = scopeFactory.CreateScope();
  102. var dbCtx = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  103. var sidMap = await dbCtx.Channel.AsNoTracking()
  104. .Where(c => c.YouTubeChannelID != null && ytChannelIds.Contains(c.YouTubeChannelID))
  105. .Select(c => new { c.YouTubeChannelID, c.SID })
  106. .ToDictionaryAsync(x => x.YouTubeChannelID!, x => x.SID, ct);
  107. var now = DateTime.UtcNow;
  108. var updated = 0;
  109. var ended = 0;
  110. foreach (var info in liveList)
  111. {
  112. if (ct.IsCancellationRequested)
  113. {
  114. break;
  115. }
  116. if (!viewers.TryGetValue(info.VideoId, out var cv))
  117. {
  118. // concurrentViewers 미제공 = 라이브 종료로 판단 (Redis 정리 + 브로드캐스트)
  119. // 만약 일시적 0명 상태였다면 다음 PubSub/FeedPoll 사이클에서 다시 마킹됨
  120. await liveStateStore.ClearLiveAsync(info.ChannelId);
  121. if (sidMap.TryGetValue(info.ChannelId, out var endedSid))
  122. {
  123. await statusBroadcaster.BroadcastAsync(endedSid, isLive: false, viewerCount: 0, videoId: null, ct);
  124. }
  125. ended++;
  126. continue;
  127. }
  128. await liveStateStore.UpdateViewerCountAsync(info.ChannelId, cv, now);
  129. if (sidMap.TryGetValue(info.ChannelId, out var sid))
  130. {
  131. await statusBroadcaster.BroadcastAsync(sid, isLive: true, viewerCount: cv, videoId: info.VideoId, ct);
  132. }
  133. updated++;
  134. }
  135. logger.LogInformation("[ViewerPoller] 완료 — live={Live}, updated={Updated}, ended={Ended}",
  136. liveList.Count, updated, ended);
  137. }
  138. finally
  139. {
  140. // 내 인스턴스가 잡은 락만 해제 (만료된 락을 다른 인스턴스가 잡았을 경우 보호)
  141. var current = await db.StringGetAsync(LockKey);
  142. if (current == _instanceId)
  143. {
  144. await db.KeyDeleteAsync(LockKey);
  145. }
  146. }
  147. }
  148. }