RedisStockBoardTrendingCache.cs 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107
  1. using Application.Abstractions.Data;
  2. using Application.Abstractions.Stocks;
  3. using Microsoft.EntityFrameworkCore;
  4. using StackExchange.Redis;
  5. namespace Infrastructure.Stocks;
  6. /// <summary>
  7. /// Redis Sorted Set 기반 종목 트렌딩 캐시 (d2 §④·§⑥) — RedisChatLeaderboard/PaperLeaderboardCache 패턴 복제.
  8. /// 점수 = 최근 24h PostStockTag 언급량 × RecentWeight + StockBoardStats.Posts(총 글 수, tie-breaker).
  9. /// TTL 10분 — 만료 시 다음 호출이 DB 집계 1회 후 재시드. StockMentionAggregator 배치가 없을 때도 lazy 로 동작.
  10. /// </summary>
  11. internal sealed class RedisStockBoardTrendingCache(IConnectionMultiplexer redis) : IStockBoardTrendingCache
  12. {
  13. private static readonly TimeSpan CacheTtl = TimeSpan.FromMinutes(10);
  14. // 최근성 가중 — 24h 언급 1건이 총 글 10건과 동일 가중.
  15. private const int RecentWeight = 10;
  16. private const string Key = "stock-trending";
  17. public async Task<IReadOnlyList<StockTrendingEntry>> GetTopAsync(int count, IAppDbContext db, CancellationToken ct = default)
  18. {
  19. if (count <= 0)
  20. {
  21. return [];
  22. }
  23. var r = redis.GetDatabase();
  24. await EnsureSeededAsync(r, db, ct);
  25. var entries = await r.SortedSetRangeByRankWithScoresAsync(Key, 0, count - 1, Order.Descending);
  26. var rows = new List<StockTrendingEntry>(entries.Length);
  27. foreach (var entry in entries)
  28. {
  29. var code = (string?)entry.Element;
  30. if (!string.IsNullOrEmpty(code))
  31. {
  32. rows.Add(new StockTrendingEntry(code, (int)entry.Score));
  33. }
  34. }
  35. return rows;
  36. }
  37. public async Task InvalidateAsync(CancellationToken ct = default)
  38. {
  39. var r = redis.GetDatabase();
  40. await r.KeyDeleteAsync(Key);
  41. }
  42. /// <summary>캐시 미스 → DB 집계 후 ZADD 재시드. TTL 동안 후속 호출은 Redis only.</summary>
  43. private static async Task EnsureSeededAsync(IDatabase r, IAppDbContext db, CancellationToken ct)
  44. {
  45. if (await r.KeyExistsAsync(Key))
  46. {
  47. return;
  48. }
  49. var since = DateTime.UtcNow.AddHours(-24);
  50. // 최근 24h 언급량 (코드별 count)
  51. var recentMentions = await db.PostStockTag.AsNoTracking()
  52. .Where(c => c.CreatedAt >= since)
  53. .GroupBy(c => c.StockCode)
  54. .Select(g => new { Code = g.Key, Count = g.Count() })
  55. .ToListAsync(ct);
  56. // 총 글 수 (StockBoardStats 파티션)
  57. var totals = await db.StockBoardStats.AsNoTracking()
  58. .Select(c => new { c.StockCode, c.Posts })
  59. .ToListAsync(ct);
  60. var scores = new Dictionary<string, int>(StringComparer.Ordinal);
  61. foreach (var t in totals)
  62. {
  63. scores[t.StockCode] = t.Posts;
  64. }
  65. foreach (var m in recentMentions)
  66. {
  67. scores.TryGetValue(m.Code, out var baseScore);
  68. scores[m.Code] = baseScore + (m.Count * RecentWeight);
  69. }
  70. if (scores.Count == 0)
  71. {
  72. return;
  73. }
  74. var entries = scores
  75. .Where(kv => kv.Value > 0)
  76. .Select(kv => new SortedSetEntry(kv.Key, kv.Value))
  77. .ToArray();
  78. if (entries.Length == 0)
  79. {
  80. return;
  81. }
  82. await r.SortedSetAddAsync(Key, entries);
  83. await r.KeyExpireAsync(Key, CacheTtl);
  84. }
  85. }