DisclosureSyncService.cs 10.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217
  1. using Application.Abstractions.Data;
  2. using Application.Helpers;
  3. using Domain.Entities.Stocks;
  4. using Microsoft.EntityFrameworkCore;
  5. using Microsoft.Extensions.DependencyInjection;
  6. using Microsoft.Extensions.Logging;
  7. using Microsoft.Extensions.Options;
  8. using SharedKernel;
  9. namespace Infrastructure.StockData;
  10. /// <summary>
  11. /// 전자공시(DART) 공시 목록 수집 (OpenDART list.json) — 기본 08:30 KST 실행. 접수일 [today-BackfillDays, today] 범위를
  12. /// page_count=100 으로 페이징(total_page 만큼 순회)하며 수집한다. KRX 배치처럼 날짜별 반복이 아니라 DART 는 날짜범위 페이징이므로
  13. /// KrxBackfill 대신 자체 페이징 루프를 쓴다. 접수번호(rcept_no) 로 upsert 하되 이미 있으면 skip(idempotent) — 갱신 없이 신규만 삽입.
  14. /// • 스펙 클램프: corp_code 미지정 검색은 최근 3개월만 제공 — BackfillDays 가 90 초과여도 창을 3개월로 클램프(경고).
  15. /// • 초기 백필(창의 가장 오래된 주에 적재분 0건): bgn_de/end_de 를 주 단위로 분할 순회해 5만 행(MaxPages×100) 상한을 회피.
  16. /// • 증분 수집: 한 페이지 전체가 기존 건이면(1페이지 제외) 그보다 오래된 페이지도 기수집으로 보고 조기 종료 — 매일 전량 재페이징 방지.
  17. /// • 페이지 상한 도달 시 잔여 페이지 수를 경고 (정렬 desc 라 오래된 공시부터 유실).
  18. /// status "013"(데이터 없음) 은 정상 종료, "000" 외 다른 status 는 오류로 보고 후 중단(재시도 대상 — DescribeStatus 라벨 로깅).
  19. /// ApiKey 미설정 시 로그만 남기고 skip. list.json 이 stock_code 를 포함하므로 corpCode.xml 매핑 불필요.
  20. /// </summary>
  21. internal sealed class DisclosureSyncService(
  22. IServiceScopeFactory scopeFactory,
  23. IHttpClientFactory httpClientFactory,
  24. IOptions<AppSettings> settings,
  25. ILogger<DisclosureSyncService> logger
  26. ) : DailyScheduledService(logger)
  27. {
  28. private const int PageCount = 100;
  29. private const int MaxPages = 500; // 안전 상한 (total_page 무한루프 방지)
  30. protected override string JobName => "DisclosureSync";
  31. protected override TimeOnly TargetTime => ParseTime(settings.Value.OpenDart.DisclosureSyncTime, new TimeOnly(8, 30));
  32. protected override int MaxRetryCount => 2;
  33. protected override TimeSpan RetryDelay => TimeSpan.FromHours(1);
  34. protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
  35. {
  36. var cfg = settings.Value.OpenDart;
  37. using var scope = scopeFactory.CreateScope();
  38. var collectorSettings = scope.ServiceProvider.GetRequiredService<ICollectorSettingsProvider>();
  39. if (!await collectorSettings.IsEnabledAsync(CollectorFlag.OpenDartDisclosure, ct))
  40. {
  41. return true;
  42. }
  43. cfg = cfg with { ApiKey = await collectorSettings.GetKeyAsync(CollectorKey.OpenDart, ct) ?? cfg.ApiKey };
  44. if (string.IsNullOrWhiteSpace(cfg.ApiKey))
  45. {
  46. Logger.LogWarning("[{Job}] OpenDart:ApiKey 미설정 — 수집 skip", JobName);
  47. return true;
  48. }
  49. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  50. var client = httpClientFactory.CreateClient(OpenDartHttp.ClientName);
  51. var days = cfg.BackfillDays > 0 ? cfg.BackfillDays : 30;
  52. var beginDate = todayKst.AddDays(-days);
  53. // corp_code 미지정 list.json 은 스펙상 최근 3개월만 제공 — 초과 설정 시 클램프
  54. var (clampedBegin, clamped) = DartBackfill.ClampToSpecWindow(beginDate, todayKst);
  55. if (clamped)
  56. {
  57. Logger.LogWarning("[{Job}] BackfillDays={Days} 가 스펙 한도(최근 {Months}개월)를 초과 — 창 시작일을 {Begin} 으로 클램프",
  58. JobName, days, DartBackfill.SpecWindowMonths, clampedBegin);
  59. beginDate = clampedBegin;
  60. }
  61. var baseUrl = cfg.BaseUrl.TrimEnd('/');
  62. // 초기 백필 판정 — 창의 가장 오래된 주(7일)에 적재된 공시가 하나도 없으면 미커버로 간주 (영업일마다 수백 건이 접수되므로 안전)
  63. var oldestWeekEnd = beginDate.AddDays(7);
  64. var isInitialBackfill = !await db.Disclosure.AsNoTracking().AnyAsync(c => c.RceptDt >= beginDate && c.RceptDt < oldestWeekEnd, ct);
  65. var inserted = 0;
  66. var skipped = 0;
  67. var ok = true;
  68. if (isInitialBackfill)
  69. {
  70. // 최초 적재: 창 전체가 5만 행(MaxPages×PageCount) 상한을 넘을 수 있어 주 단위로 분할 순회 (사업보고서 시즌 30일 창이 5만 초과 현실적)
  71. var chunks = DartBackfill.SplitWeekly(beginDate, todayKst);
  72. Logger.LogInformation("[{Job}] 초기 백필 — 범위=[{Begin}~{Today}] 를 주 단위 {Chunks}개 창으로 분할 수집", JobName, beginDate, todayKst, chunks.Count);
  73. foreach (var (chunkBegin, chunkEnd) in chunks)
  74. {
  75. var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, chunkBegin, chunkEnd, allowEarlyStop: false, ct);
  76. inserted += result.Inserted;
  77. skipped += result.Skipped;
  78. if (!result.Ok)
  79. {
  80. ok = false;
  81. break;
  82. }
  83. await Task.Delay(300, ct); // quota 보호 (청크 간)
  84. }
  85. }
  86. else
  87. {
  88. // 증분 수집: 신규분이 최신 페이지에 몰려 있으므로 전체 기존 페이지를 만나면 조기 종료
  89. var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, beginDate, todayKst, allowEarlyStop: true, ct);
  90. inserted = result.Inserted;
  91. skipped = result.Skipped;
  92. ok = result.Ok;
  93. }
  94. Logger.LogInformation("[{Job}] 완료 — 범위=[{Begin}~{Today}], inserted={Inserted}, skipped={Skipped}, 초기백필={Initial}",
  95. JobName, beginDate, todayKst, inserted, skipped, isInitialBackfill);
  96. return ok;
  97. }
  98. /// <summary>
  99. /// [bgnDe, endDe] 접수일 범위를 desc 정렬로 페이징 수집. 신규만 삽입(rcept_no idempotent).
  100. /// allowEarlyStop=true 면 "페이지 전체 기존 + 1페이지 아님" 조건에서 조기 종료 (DartBackfill.ShouldStopEarly).
  101. /// API 오류 status 는 Ok=false 로 반환해 베이스 재시도로 연결한다.
  102. /// </summary>
  103. private async Task<(bool Ok, int Inserted, int Skipped)> SyncRangeAsync(
  104. IAppDbContext db,
  105. HttpClient client,
  106. string baseUrl,
  107. string apiKey,
  108. DateOnly bgnDe,
  109. DateOnly endDe,
  110. bool allowEarlyStop,
  111. CancellationToken ct)
  112. {
  113. var inserted = 0;
  114. var skipped = 0;
  115. var page = 1;
  116. var totalPage = 1;
  117. while (page <= totalPage && page <= MaxPages)
  118. {
  119. ct.ThrowIfCancellationRequested();
  120. var url = $"{baseUrl}/api/list.json?crtfc_key={apiKey}&bgn_de={bgnDe:yyyyMMdd}&end_de={endDe:yyyyMMdd}&page_no={page}&page_count={PageCount}";
  121. var json = await OpenDartHttp.GetStringWithRetryAsync(client, url, Logger, ct);
  122. var result = DartDisclosureParser.Parse(json);
  123. if (result.Status == "013")
  124. {
  125. Logger.LogInformation("[{Job}] 조회된 데이터 없음(013) — 종료. 범위=[{Begin}~{End}]", JobName, bgnDe, endDe);
  126. break;
  127. }
  128. if (result.Status != "000")
  129. {
  130. Logger.LogError("[{Job}] OpenDART 오류 status={Status}({Label}) message={Message} (page={Page}, 범위=[{Begin}~{End}])",
  131. JobName, result.Status, DartDisclosureParser.DescribeStatus(result.Status), result.Message, page, bgnDe, endDe);
  132. return (false, inserted, skipped);
  133. }
  134. totalPage = result.TotalPage > 0 ? result.TotalPage : 1;
  135. if (result.Rows.Count == 0)
  136. {
  137. break;
  138. }
  139. var rceptNos = result.Rows.Select(c => c.RceptNo).ToList();
  140. var existing = (await db.Disclosure.AsNoTracking()
  141. .Where(c => rceptNos.Contains(c.RceptNo))
  142. .Select(c => c.RceptNo)
  143. .ToListAsync(ct)).ToHashSet();
  144. var newCount = 0;
  145. foreach (var row in result.Rows)
  146. {
  147. if (existing.Contains(row.RceptNo))
  148. {
  149. skipped++;
  150. continue;
  151. }
  152. var created = Disclosure.Create(row.RceptNo, row.CorpCode, row.CorpName, row.StockCode, row.CorpCls, row.ReportNm, row.FlrNm, row.RceptDt, row.Rm);
  153. await db.Disclosure.AddAsync(created, ct);
  154. existing.Add(row.RceptNo); // 같은 실행 내 중복 접수번호 방어
  155. inserted++;
  156. newCount++;
  157. }
  158. await db.SaveChangesAsync(ct);
  159. Logger.LogInformation("[{Job}] page {Page}/{TotalPage} 처리 — rows={Rows}, 신규={New} (누적 inserted={Inserted}, skipped={Skipped})",
  160. JobName, page, totalPage, result.Rows.Count, newCount, inserted, skipped);
  161. if (allowEarlyStop && DartBackfill.ShouldStopEarly(page, result.Rows.Count, newCount))
  162. {
  163. Logger.LogInformation("[{Job}] 조기 종료 — page {Page} 전체가 기수집 건 (이후 페이지는 더 오래된 기수집분)", JobName, page);
  164. return (true, inserted, skipped);
  165. }
  166. page++;
  167. if (page <= totalPage)
  168. {
  169. await Task.Delay(300, ct); // quota 보호
  170. }
  171. }
  172. if (totalPage > MaxPages)
  173. {
  174. Logger.LogWarning("[{Job}] 페이지 상한 도달 — totalPage={TotalPage}, 상한={MaxPages}, 미수집 잔여 {Remaining}페이지 (정렬 desc 라 가장 오래된 공시부터 유실). 범위=[{Begin}~{End}] 를 좁혀 재수집 필요",
  175. JobName, totalPage, MaxPages, totalPage - MaxPages, bgnDe, endDe);
  176. }
  177. return (true, inserted, skipped);
  178. }
  179. }