StockMasterSyncService.cs 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186
  1. using Application.Abstractions.Data;
  2. using Application.Helpers;
  3. using Domain.Entities.Stocks;
  4. using Domain.Entities.Stocks.ValueObject;
  5. using Microsoft.EntityFrameworkCore;
  6. using Microsoft.Extensions.DependencyInjection;
  7. using Microsoft.Extensions.Logging;
  8. using Microsoft.Extensions.Options;
  9. using SharedKernel;
  10. namespace Infrastructure.StockData;
  11. /// <summary>
  12. /// 종목 마스터 동기화 — 금융위 KRX 상장종목정보 API 를 일 1회(기본 07:30 KST) 전량 수집하여 Stock upsert.
  13. /// 신규 상장 = insert, 명칭/시장 변경 = update, 스냅샷에서 사라진 종목 = 상폐 soft-off.
  14. /// 단일 엔드포인트 페이징 수집이라 조기 종료/페이지 상한 시 부분 스냅샷이 될 수 있으므로,
  15. /// 상폐 sweep 은 시장별 가드(StockMasterDelist — 0건 또는 활성 대비 80% 미만이면 보류+경고)를 통과한 시장에서만 수행.
  16. /// ServiceKey 미설정 시 로그만 남기고 skip.
  17. /// </summary>
  18. internal sealed class StockMasterSyncService(
  19. IServiceScopeFactory scopeFactory,
  20. IHttpClientFactory httpClientFactory,
  21. IOptions<AppSettings> settings,
  22. ILogger<StockMasterSyncService> logger
  23. ) : DailyScheduledService(logger)
  24. {
  25. private const string ServicePath = "/1160100/service/GetKrxListedInfoService/getItemInfo";
  26. private const int MaxPages = 50;
  27. private const int MaxBaseDateLookback = 7;
  28. protected override string JobName => "StockMasterSync";
  29. protected override TimeOnly TargetTime => ParseTime(settings.Value.StockData.MasterSyncTime, new TimeOnly(7, 30));
  30. protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
  31. {
  32. var cfg = settings.Value.StockData.DataGoKr;
  33. using var scope = scopeFactory.CreateScope();
  34. var collectorSettings = scope.ServiceProvider.GetRequiredService<ICollectorSettingsProvider>();
  35. if (!await collectorSettings.IsEnabledAsync(CollectorFlag.StockDataMaster, ct))
  36. {
  37. return true;
  38. }
  39. cfg = cfg with { ServiceKey = await collectorSettings.GetKeyAsync(CollectorKey.DataGoKr, ct) ?? cfg.ServiceKey };
  40. if (string.IsNullOrWhiteSpace(cfg.ServiceKey))
  41. {
  42. Logger.LogWarning("[{Job}] StockData:DataGoKr:ServiceKey 미설정 — 수집 skip", JobName);
  43. return true;
  44. }
  45. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  46. var client = httpClientFactory.CreateClient(DataGoKrHttp.ClientName);
  47. // 상장종목정보는 basDt 단위 스냅샷 — 직전 영업일부터 최대 7일 소급하며 데이터가 있는 기준일을 찾는다
  48. var baseDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, todayKst.AddDays(1), ct);
  49. List<DataGoKrStockParser.ListedItem>? snapshot = null;
  50. for (var back = 0; back < MaxBaseDateLookback; back++)
  51. {
  52. snapshot = await FetchSnapshotAsync(client, cfg, baseDate, ct);
  53. if (snapshot.Count > 0)
  54. {
  55. break;
  56. }
  57. baseDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, baseDate, ct);
  58. }
  59. if (snapshot is null || snapshot.Count == 0)
  60. {
  61. Logger.LogError("[{Job}] 상장종목 스냅샷 없음 — 최근 {Days}영업일 조회 실패", JobName, MaxBaseDateLookback);
  62. return true;
  63. }
  64. // 동일 코드 중복 행은 마지막 행 우선. 상폐 sweep 가드용으로 시장별 건수도 함께 집계 (0건 시장 포함)
  65. var byCode = new Dictionary<string, DataGoKrStockParser.ListedItem>();
  66. var snapshotCountByMarket = new Dictionary<StockMarket, int> { [StockMarket.KOSPI] = 0, [StockMarket.KOSDAQ] = 0, [StockMarket.KONEX] = 0 };
  67. foreach (var item in snapshot)
  68. {
  69. var market = MapMarket(item.MarketName);
  70. if (market is null)
  71. {
  72. continue;
  73. }
  74. if (!byCode.ContainsKey(item.Code))
  75. {
  76. snapshotCountByMarket[market.Value]++;
  77. }
  78. byCode[item.Code] = item;
  79. }
  80. var stocks = await db.Stock.ToListAsync(ct);
  81. var stockByCode = stocks.ToDictionary(c => c.Code);
  82. var inserted = 0;
  83. var updated = 0;
  84. var delisted = 0;
  85. foreach (var (code, item) in byCode)
  86. {
  87. var market = MapMarket(item.MarketName)!.Value;
  88. if (stockByCode.TryGetValue(code, out var stock))
  89. {
  90. var wasUpdated = stock.UpdatedAt;
  91. stock.UpdateMaster(item.Name, market, item.Isin);
  92. if (stock.UpdatedAt != wasUpdated)
  93. {
  94. updated++;
  95. }
  96. }
  97. else
  98. {
  99. await db.Stock.AddAsync(Stock.Create(code, item.Name, market, item.BaseDate, item.Isin), ct);
  100. inserted++;
  101. }
  102. }
  103. // 스냅샷에서 사라진 활성 종목 = 상폐 soft-off — 단, 시장별 가드(빈 응답/활성 대비 80% 미만)를 통과한 시장만
  104. var activeCountByMarket = stocks.Where(c => c.IsActive).GroupBy(c => c.Market).ToDictionary(c => c.Key, c => c.Count());
  105. var delistableMarkets = StockMasterDelist.GetDelistableMarkets(snapshotCountByMarket, activeCountByMarket).ToHashSet();
  106. foreach (var (market, snapshotCount) in snapshotCountByMarket.Where(c => !delistableMarkets.Contains(c.Key)))
  107. {
  108. Logger.LogWarning("[{Job}] {Market} 상폐 sweep 보류 — snapshot={Snapshot}, active={Active} (빈/부분 응답 방어)",
  109. JobName, market, snapshotCount, activeCountByMarket.TryGetValue(market, out var active) ? active : 0);
  110. }
  111. foreach (var stock in stocks.Where(c => c.IsActive && delistableMarkets.Contains(c.Market) && !byCode.ContainsKey(c.Code)))
  112. {
  113. stock.MarkDelisted(baseDate);
  114. delisted++;
  115. }
  116. await db.SaveChangesAsync(ct);
  117. Logger.LogInformation("[{Job}] 완료 — basDt={BaseDate}, snapshot={Snapshot}, inserted={Inserted}, updated={Updated}, delisted={Delisted}",
  118. JobName, baseDate, byCode.Count, inserted, updated, delisted);
  119. return true;
  120. }
  121. private async Task<List<DataGoKrStockParser.ListedItem>> FetchSnapshotAsync(HttpClient client, AppSettings.StockDataSection.DataGoKrSection cfg, DateOnly baseDate, CancellationToken ct)
  122. {
  123. var all = new List<DataGoKrStockParser.ListedItem>();
  124. var totalCount = int.MaxValue;
  125. for (var pageNo = 1; pageNo <= MaxPages && all.Count < totalCount; pageNo++)
  126. {
  127. var url = $"{cfg.BaseUrl.TrimEnd('/')}{ServicePath}?serviceKey={Uri.EscapeDataString(cfg.ServiceKey)}&resultType=json&numOfRows={cfg.PageSize}&pageNo={pageNo}&basDt={baseDate:yyyyMMdd}";
  128. var json = await DataGoKrHttp.GetStringWithRetryAsync(client, url, Logger, ct);
  129. var (items, total) = DataGoKrStockParser.ParseListedItems(json);
  130. totalCount = total;
  131. if (items.Count == 0)
  132. {
  133. break;
  134. }
  135. all.AddRange(items);
  136. }
  137. return all;
  138. }
  139. private static StockMarket? MapMarket(string marketName)
  140. {
  141. if (marketName.Contains("KOSPI", StringComparison.OrdinalIgnoreCase))
  142. {
  143. return StockMarket.KOSPI;
  144. }
  145. if (marketName.Contains("KOSDAQ", StringComparison.OrdinalIgnoreCase))
  146. {
  147. return StockMarket.KOSDAQ;
  148. }
  149. if (marketName.Contains("KONEX", StringComparison.OrdinalIgnoreCase))
  150. {
  151. return StockMarket.KONEX;
  152. }
  153. return null;
  154. }
  155. }