KrxStockMasterSyncService.cs 7.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157
  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 OpenAPI) — stk_isu_base_info(KOSPI) + ksq_isu_base_info(KOSDAQ) + knx_isu_base_info(KONEX) 종목기본정보를
  13. /// 일 1회(기본 07:40 KST) 전량 수집하여 Stock upsert. 신규 상장 = insert, 명칭/시장/영문명/업종 변경 = update,
  14. /// 스냅샷에서 사라진 활성 종목 = 상폐 soft-off. ApiKey 미설정 시 로그만 남기고 skip (data.go.kr 배치와 동일 정책).
  15. /// 종목기본정보는 basDd 단위 스냅샷 — 직전 영업일부터 최대 7일 소급하며 데이터가 있는 기준일을 찾는다.
  16. /// 상폐 sweep 은 시장별 가드(StockMasterDelist)를 통과한 시장에서만 수행 — 한 시장만 빈/부분 응답일 때
  17. /// 그 시장 전체(예: KOSDAQ ~1,700종목)가 하루 동안 대량 상폐되던 결함 방지 (0건 또는 활성 대비 80% 미만이면 보류+경고).
  18. /// KONEX 응답 필드셋은 KOSPI/KOSDAQ 과 동일(ISU_SRT_CD/ISU_CD/ISU_ABBRV/ISU_ENG_NM/LIST_DD/MKT_TP_NM/SECT_TP_NM) — KrxStockParser 재사용.
  19. /// </summary>
  20. internal sealed class KrxStockMasterSyncService(
  21. IServiceScopeFactory scopeFactory,
  22. IHttpClientFactory httpClientFactory,
  23. IOptions<AppSettings> settings,
  24. ILogger<KrxStockMasterSyncService> logger
  25. ) : DailyScheduledService(logger)
  26. {
  27. // (시장, 엔드포인트 경로) — KOSPI/KOSDAQ/KONEX. 세 시장 모두 Stock 테이블에 Market 으로만 구분되어 적재된다.
  28. private static readonly (StockMarket Market, string Path)[] Endpoints =
  29. [
  30. (StockMarket.KOSPI, "/svc/apis/sto/stk_isu_base_info"),
  31. (StockMarket.KOSDAQ, "/svc/apis/sto/ksq_isu_base_info"),
  32. (StockMarket.KONEX, "/svc/apis/sto/knx_isu_base_info")
  33. ];
  34. private const int MaxBaseDateLookback = 7;
  35. protected override string JobName => "KrxStockMasterSync";
  36. protected override TimeOnly TargetTime => ParseTime(settings.Value.KRXCoKr.MasterSyncTime, new TimeOnly(7, 40));
  37. protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
  38. {
  39. var cfg = settings.Value.KRXCoKr;
  40. using var scope = scopeFactory.CreateScope();
  41. var collectorSettings = scope.ServiceProvider.GetRequiredService<ICollectorSettingsProvider>();
  42. if (!await collectorSettings.IsEnabledAsync(CollectorFlag.KrxStock, ct))
  43. {
  44. return true;
  45. }
  46. cfg = cfg with { ApiKey = await collectorSettings.GetKeyAsync(CollectorKey.Krx, ct) ?? cfg.ApiKey };
  47. if (string.IsNullOrWhiteSpace(cfg.ApiKey))
  48. {
  49. Logger.LogWarning("[{Job}] KRXCoKr:ApiKey 미설정 — 수집 skip", JobName);
  50. return true;
  51. }
  52. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  53. var client = httpClientFactory.CreateClient(KrxCoKrHttp.ClientName);
  54. // 종목기본정보는 basDd 단위 스냅샷 — 직전 영업일부터 최대 7일 소급하며 데이터가 있는 기준일을 찾는다
  55. var baseDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, todayKst.AddDays(1), ct);
  56. Dictionary<StockMarket, List<KrxStockParser.MasterItem>>? snapshot = null;
  57. for (var back = 0; back < MaxBaseDateLookback; back++)
  58. {
  59. snapshot = await FetchSnapshotAsync(client, cfg, baseDate, ct);
  60. if (snapshot.Values.Sum(c => c.Count) > 0)
  61. {
  62. break;
  63. }
  64. baseDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, baseDate, ct);
  65. }
  66. if (snapshot is null || snapshot.Values.Sum(c => c.Count) == 0)
  67. {
  68. Logger.LogError("[{Job}] 종목기본정보 스냅샷 없음 — 최근 {Days}영업일 조회 실패", JobName, MaxBaseDateLookback);
  69. return true;
  70. }
  71. // 동일 코드 중복 행은 마지막 행 우선
  72. var byCode = new Dictionary<string, KrxStockParser.MasterItem>();
  73. foreach (var item in snapshot.Values.SelectMany(c => c))
  74. {
  75. byCode[item.Code] = item;
  76. }
  77. var stocks = await db.Stock.ToListAsync(ct);
  78. var stockByCode = stocks.ToDictionary(c => c.Code);
  79. var inserted = 0;
  80. var updated = 0;
  81. var delisted = 0;
  82. foreach (var (code, item) in byCode)
  83. {
  84. if (stockByCode.TryGetValue(code, out var stock))
  85. {
  86. var wasUpdated = stock.UpdatedAt;
  87. stock.UpdateMaster(item.Name, item.Market, item.Isin, item.EnglishName, item.SectorName);
  88. if (stock.UpdatedAt != wasUpdated)
  89. {
  90. updated++;
  91. }
  92. }
  93. else
  94. {
  95. await db.Stock.AddAsync(Stock.Create(code, item.Name, item.Market, item.ListedDate, item.Isin, englishName: item.EnglishName, sectorName: item.SectorName), ct);
  96. inserted++;
  97. }
  98. }
  99. // 스냅샷에서 사라진 활성 종목 = 상폐 soft-off — 단, 시장별 가드(빈 응답/활성 대비 80% 미만)를 통과한 시장만
  100. var snapshotCountByMarket = snapshot.ToDictionary(c => c.Key, c => c.Value.Count);
  101. var activeCountByMarket = stocks.Where(c => c.IsActive).GroupBy(c => c.Market).ToDictionary(c => c.Key, c => c.Count());
  102. var delistableMarkets = StockMasterDelist.GetDelistableMarkets(snapshotCountByMarket, activeCountByMarket).ToHashSet();
  103. foreach (var (market, snapshotCount) in snapshotCountByMarket.Where(c => !delistableMarkets.Contains(c.Key)))
  104. {
  105. Logger.LogWarning("[{Job}] {Market} 상폐 sweep 보류 — snapshot={Snapshot}, active={Active} (빈/부분 응답 방어)",
  106. JobName, market, snapshotCount, activeCountByMarket.TryGetValue(market, out var active) ? active : 0);
  107. }
  108. foreach (var stock in stocks.Where(c => c.IsActive && delistableMarkets.Contains(c.Market) && !byCode.ContainsKey(c.Code)))
  109. {
  110. stock.MarkDelisted(baseDate);
  111. delisted++;
  112. }
  113. await db.SaveChangesAsync(ct);
  114. Logger.LogInformation("[{Job}] 완료 — basDd={BaseDate}, snapshot={Snapshot}, inserted={Inserted}, updated={Updated}, delisted={Delisted}",
  115. JobName, baseDate, byCode.Count, inserted, updated, delisted);
  116. return true;
  117. }
  118. /// <summary>시장별 스냅샷 수집 — 상폐 sweep 가드를 위해 시장 단위 건수를 보존해 반환한다 (0건 시장도 키 존재).</summary>
  119. private async Task<Dictionary<StockMarket, List<KrxStockParser.MasterItem>>> FetchSnapshotAsync(HttpClient client, AppSettings.KRXCoKrSection cfg, DateOnly baseDate, CancellationToken ct)
  120. {
  121. var byMarket = new Dictionary<StockMarket, List<KrxStockParser.MasterItem>>();
  122. foreach (var (market, path) in Endpoints)
  123. {
  124. var url = $"{cfg.BaseUrl.TrimEnd('/')}{path}?basDd={baseDate:yyyyMMdd}";
  125. var json = await KrxCoKrHttp.GetStringWithRetryAsync(client, url, cfg.ApiKey, Logger, ct);
  126. var rows = KrxStockParser.ParseMasterInfo(json, market);
  127. Logger.LogInformation("[{Job}] {Market} basDd={BaseDate} rows={Rows}", JobName, market, baseDate, rows.Count);
  128. byMarket[market] = rows.ToList();
  129. }
  130. return byMarket;
  131. }
  132. }