RankingDailyAggregatorService.cs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315
  1. using Application.Abstractions.Data;
  2. using Domain.Entities.Donations;
  3. using Domain.Entities.Donations.ValueObject;
  4. using Microsoft.EntityFrameworkCore;
  5. using Microsoft.Extensions.DependencyInjection;
  6. using Microsoft.Extensions.Hosting;
  7. using Microsoft.Extensions.Logging;
  8. namespace Infrastructure.Ranking;
  9. /// <summary>
  10. /// 매일 자정 00:30(KST) 1회 실행 — 랭킹 스냅샷 집계.
  11. /// - YouTubeDailyAggregator 완료 후 실행 (00:00 + 30분 버퍼)
  12. /// - 3가지 랭킹 × 6가지 기간 × (전체 + 카테고리) = N개의 스냅샷 생성
  13. /// - 기존 스냅샷은 삭제 후 재생성 (upsert 대신 단순 truncate+insert)
  14. ///
  15. /// 종합 점수 가중치 (Comprehensive):
  16. /// 머니 × 0.55 + 후원자수(×10000 환산) × 0.15 + 방송시간(초 / 3600) × 0.15
  17. /// + 조회수 × 0.1 + 좋아요 × 0.05
  18. /// </summary>
  19. internal sealed class RankingDailyAggregatorService(
  20. IServiceScopeFactory scopeFactory,
  21. ILogger<RankingDailyAggregatorService> logger
  22. ) : BackgroundService
  23. {
  24. private static readonly TimeSpan InitialDelay = TimeSpan.FromMinutes(3);
  25. private static readonly TimeZoneInfo Kst = TimeZoneInfo.FindSystemTimeZoneById("Korea Standard Time");
  26. private const int OffsetMinutesAfterMidnight = 30;
  27. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  28. {
  29. await Task.Delay(InitialDelay, stoppingToken);
  30. logger.LogInformation("[RankingAggregator] 서비스 시작 — 매일 KST 00:{Offset} 실행", OffsetMinutesAfterMidnight.ToString("D2"));
  31. while (!stoppingToken.IsCancellationRequested)
  32. {
  33. var delay = CalculateDelayUntilNextRun();
  34. logger.LogInformation("[RankingAggregator] 다음 실행까지 대기: {Hours}시간 {Minutes}분", delay.Hours, delay.Minutes);
  35. try
  36. {
  37. await Task.Delay(delay, stoppingToken);
  38. }
  39. catch (TaskCanceledException)
  40. {
  41. break;
  42. }
  43. try
  44. {
  45. await AggregateAllAsync(stoppingToken);
  46. }
  47. catch (Exception ex)
  48. {
  49. logger.LogError(ex, "[RankingAggregator] 집계 실행 중 오류");
  50. }
  51. }
  52. }
  53. private static TimeSpan CalculateDelayUntilNextRun()
  54. {
  55. var nowKst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, Kst);
  56. var nextRunKst = nowKst.Date.AddDays(1).AddMinutes(OffsetMinutesAfterMidnight);
  57. var delay = nextRunKst - nowKst;
  58. if (delay <= TimeSpan.Zero)
  59. {
  60. delay = TimeSpan.FromMinutes(1);
  61. }
  62. return delay;
  63. }
  64. private async Task AggregateAllAsync(CancellationToken ct)
  65. {
  66. logger.LogInformation("[RankingAggregator] 스냅샷 집계 시작");
  67. var periods = Enum.GetValues<RankingPeriod>();
  68. var types = Enum.GetValues<RankingType>();
  69. foreach (var period in periods)
  70. {
  71. foreach (var type in types)
  72. {
  73. try
  74. {
  75. await AggregateOneAsync(type, period, null, ct);
  76. }
  77. catch (Exception ex)
  78. {
  79. logger.LogError(ex, "[RankingAggregator] {Type}/{Period} 집계 실패", type, period);
  80. }
  81. }
  82. }
  83. logger.LogInformation("[RankingAggregator] 스냅샷 집계 완료");
  84. }
  85. private async Task AggregateOneAsync(RankingType type, RankingPeriod period, string? category, CancellationToken ct)
  86. {
  87. using var scope = scopeFactory.CreateScope();
  88. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  89. var (periodStart, periodEnd) = ResolvePeriodRange(period);
  90. // 기존 스냅샷 삭제 (동일 type/period/category)
  91. var existing = await db.RankingSnapshot
  92. .Where(r => r.Type == type && r.Period == period && r.Category == category)
  93. .ToListAsync(ct);
  94. if (existing.Count > 0)
  95. {
  96. db.RankingSnapshot.RemoveRange(existing);
  97. await db.SaveChangesAsync(ct);
  98. }
  99. // 전일 스냅샷 (PreviousRank 계산용)
  100. var previousRanks = existing
  101. .GroupBy(r => r.MemberID)
  102. .ToDictionary(g => g.Key, g => g.First().Rank);
  103. List<AggregatedRow> rows = type switch
  104. {
  105. RankingType.Comprehensive => await AggregateComprehensiveAsync(db, periodStart, periodEnd, ct),
  106. RankingType.Creator => await AggregateCreatorAsync(db, periodStart, periodEnd, ct),
  107. RankingType.Donor => await AggregateDonorAsync(db, periodStart, periodEnd, ct),
  108. _ => []
  109. };
  110. if (rows.Count == 0)
  111. {
  112. return;
  113. }
  114. var ranked = rows
  115. .OrderByDescending(r => r.Score)
  116. .ThenByDescending(r => r.MoneyAmount)
  117. .ToList();
  118. var now = DateTime.UtcNow;
  119. var snapshots = new List<RankingSnapshot>(ranked.Count);
  120. for (var i = 0; i < ranked.Count; i++)
  121. {
  122. var row = ranked[i];
  123. var rank = i + 1;
  124. int? previousRank = previousRanks.GetValueOrDefault(row.MemberID) is var p && p > 0 ? p : null;
  125. snapshots.Add(RankingSnapshot.Create(
  126. type,
  127. period,
  128. periodStart,
  129. periodEnd,
  130. category,
  131. row.MemberID,
  132. rank,
  133. row.Score,
  134. row.MoneyAmount,
  135. row.DonorCount,
  136. row.BroadcastSec,
  137. row.ViewCount,
  138. row.LikeCount,
  139. previousRank
  140. ));
  141. }
  142. await db.RankingSnapshot.AddRangeAsync(snapshots, ct);
  143. await db.SaveChangesAsync(ct);
  144. logger.LogInformation("[RankingAggregator] {Type}/{Period} — {Count}개 스냅샷 생성", type, period, snapshots.Count);
  145. }
  146. private static (DateTime Start, DateTime End) ResolvePeriodRange(RankingPeriod period)
  147. {
  148. var nowKst = TimeZoneInfo.ConvertTimeFromUtc(DateTime.UtcNow, Kst);
  149. var today = nowKst.Date;
  150. return period switch
  151. {
  152. RankingPeriod.Today => (ToUtc(today), ToUtc(today.AddDays(1))),
  153. RankingPeriod.Yesterday => (ToUtc(today.AddDays(-1)), ToUtc(today)),
  154. RankingPeriod.Week => (ToUtc(today.AddDays(-(int)today.DayOfWeek)), ToUtc(today.AddDays(1))),
  155. RankingPeriod.Month => (ToUtc(new DateTime(today.Year, today.Month, 1)), ToUtc(today.AddDays(1))),
  156. RankingPeriod.LastMonth => (
  157. ToUtc(new DateTime(today.Year, today.Month, 1).AddMonths(-1)),
  158. ToUtc(new DateTime(today.Year, today.Month, 1))
  159. ),
  160. RankingPeriod.All => (DateTime.MinValue.ToUniversalTime(), ToUtc(today.AddDays(1))),
  161. _ => (ToUtc(new DateTime(today.Year, today.Month, 1)), ToUtc(today.AddDays(1)))
  162. };
  163. }
  164. private static DateTime ToUtc(DateTime kstDate)
  165. {
  166. var kstUnspecified = DateTime.SpecifyKind(kstDate, DateTimeKind.Unspecified);
  167. return TimeZoneInfo.ConvertTimeToUtc(kstUnspecified, Kst);
  168. }
  169. private static async Task<List<AggregatedRow>> AggregateCreatorAsync(IAppDbContext db, DateTime start, DateTime end, CancellationToken ct)
  170. {
  171. var grouped = await db.Donation.AsNoTracking()
  172. .Where(d => d.CreatedAt >= start && d.CreatedAt < end && !d.IsTest)
  173. .GroupBy(d => d.ReceiverMemberID)
  174. .Select(g => new
  175. {
  176. MemberID = g.Key,
  177. MoneyAmount = g.Sum(d => (long)d.NetAmount),
  178. DonorCount = g.Select(d => d.SponsorMemberID).Distinct().Count()
  179. })
  180. .ToListAsync(ct);
  181. return grouped.Select(g => new AggregatedRow
  182. {
  183. MemberID = g.MemberID,
  184. MoneyAmount = g.MoneyAmount,
  185. DonorCount = g.DonorCount,
  186. BroadcastSec = 0,
  187. ViewCount = 0,
  188. LikeCount = 0,
  189. Score = g.MoneyAmount
  190. }).ToList();
  191. }
  192. private static async Task<List<AggregatedRow>> AggregateDonorAsync(IAppDbContext db, DateTime start, DateTime end, CancellationToken ct)
  193. {
  194. var grouped = await db.Donation.AsNoTracking()
  195. .Where(d => d.CreatedAt >= start && d.CreatedAt < end && !d.IsTest)
  196. .GroupBy(d => d.SponsorMemberID)
  197. .Select(g => new
  198. {
  199. MemberID = g.Key,
  200. MoneyAmount = g.Sum(d => (long)d.Amount),
  201. DonorCount = g.Count()
  202. })
  203. .ToListAsync(ct);
  204. return grouped.Select(g => new AggregatedRow
  205. {
  206. MemberID = g.MemberID,
  207. MoneyAmount = g.MoneyAmount,
  208. DonorCount = g.DonorCount,
  209. BroadcastSec = 0,
  210. ViewCount = 0,
  211. LikeCount = 0,
  212. Score = g.MoneyAmount
  213. }).ToList();
  214. }
  215. private static async Task<List<AggregatedRow>> AggregateComprehensiveAsync(IAppDbContext db, DateTime start, DateTime end, CancellationToken ct)
  216. {
  217. var creatorRows = await AggregateCreatorAsync(db, start, end, ct);
  218. var creatorMap = creatorRows.ToDictionary(r => r.MemberID, r => r);
  219. // Member ID 목록 수집 (후원 받은 회원 + 방송 활동 회원)
  220. var broadcastStats = await db.ChannelBroadcastStats.AsNoTracking()
  221. .Join(db.Channel.AsNoTracking(), s => s.ChannelID, c => c.ID, (s, c) => new
  222. {
  223. c.MemberID,
  224. s.TotalDurationSec,
  225. s.CumulativeViews,
  226. s.TotalLikes
  227. })
  228. .ToListAsync(ct);
  229. foreach (var stat in broadcastStats)
  230. {
  231. if (creatorMap.TryGetValue(stat.MemberID, out var row))
  232. {
  233. row.BroadcastSec = stat.TotalDurationSec;
  234. row.ViewCount = stat.CumulativeViews;
  235. row.LikeCount = stat.TotalLikes;
  236. }
  237. else
  238. {
  239. creatorMap[stat.MemberID] = new AggregatedRow
  240. {
  241. MemberID = stat.MemberID,
  242. MoneyAmount = 0,
  243. DonorCount = 0,
  244. BroadcastSec = stat.TotalDurationSec,
  245. ViewCount = stat.CumulativeViews,
  246. LikeCount = stat.TotalLikes
  247. };
  248. }
  249. }
  250. foreach (var row in creatorMap.Values)
  251. {
  252. row.Score = CalculateComprehensiveScore(row);
  253. }
  254. return creatorMap.Values.ToList();
  255. }
  256. private static long CalculateComprehensiveScore(AggregatedRow row)
  257. {
  258. // 머니 × 0.55 + 후원자수 × 10000 × 0.15 + 방송시간/3600(시) × 10000 × 0.15 + 조회수 × 0.1 + 좋아요 × 0.05
  259. var moneyScore = row.MoneyAmount * 0.55;
  260. var donorScore = row.DonorCount * 10000.0 * 0.15;
  261. var broadcastScore = (row.BroadcastSec / 3600.0) * 10000.0 * 0.15;
  262. var viewScore = row.ViewCount * 0.1;
  263. var likeScore = row.LikeCount * 0.05;
  264. return (long)(moneyScore + donorScore + broadcastScore + viewScore + likeScore);
  265. }
  266. private sealed class AggregatedRow
  267. {
  268. public int MemberID { get; set; }
  269. public long Score { get; set; }
  270. public long MoneyAmount { get; set; }
  271. public int DonorCount { get; set; }
  272. public long BroadcastSec { get; set; }
  273. public long ViewCount { get; set; }
  274. public long LikeCount { get; set; }
  275. }
  276. }