DisclosureSyncService.cs 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210
  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. if (string.IsNullOrWhiteSpace(cfg.ApiKey))
  38. {
  39. Logger.LogWarning("[{Job}] OpenDart:ApiKey 미설정 — 수집 skip", JobName);
  40. return true;
  41. }
  42. using var scope = scopeFactory.CreateScope();
  43. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  44. var client = httpClientFactory.CreateClient(OpenDartHttp.ClientName);
  45. var days = cfg.BackfillDays > 0 ? cfg.BackfillDays : 30;
  46. var beginDate = todayKst.AddDays(-days);
  47. // corp_code 미지정 list.json 은 스펙상 최근 3개월만 제공 — 초과 설정 시 클램프
  48. var (clampedBegin, clamped) = DartBackfill.ClampToSpecWindow(beginDate, todayKst);
  49. if (clamped)
  50. {
  51. Logger.LogWarning("[{Job}] BackfillDays={Days} 가 스펙 한도(최근 {Months}개월)를 초과 — 창 시작일을 {Begin} 으로 클램프",
  52. JobName, days, DartBackfill.SpecWindowMonths, clampedBegin);
  53. beginDate = clampedBegin;
  54. }
  55. var baseUrl = cfg.BaseUrl.TrimEnd('/');
  56. // 초기 백필 판정 — 창의 가장 오래된 주(7일)에 적재된 공시가 하나도 없으면 미커버로 간주 (영업일마다 수백 건이 접수되므로 안전)
  57. var oldestWeekEnd = beginDate.AddDays(7);
  58. var isInitialBackfill = !await db.Disclosure.AsNoTracking().AnyAsync(c => c.RceptDt >= beginDate && c.RceptDt < oldestWeekEnd, ct);
  59. var inserted = 0;
  60. var skipped = 0;
  61. var ok = true;
  62. if (isInitialBackfill)
  63. {
  64. // 최초 적재: 창 전체가 5만 행(MaxPages×PageCount) 상한을 넘을 수 있어 주 단위로 분할 순회 (사업보고서 시즌 30일 창이 5만 초과 현실적)
  65. var chunks = DartBackfill.SplitWeekly(beginDate, todayKst);
  66. Logger.LogInformation("[{Job}] 초기 백필 — 범위=[{Begin}~{Today}] 를 주 단위 {Chunks}개 창으로 분할 수집", JobName, beginDate, todayKst, chunks.Count);
  67. foreach (var (chunkBegin, chunkEnd) in chunks)
  68. {
  69. var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, chunkBegin, chunkEnd, allowEarlyStop: false, ct);
  70. inserted += result.Inserted;
  71. skipped += result.Skipped;
  72. if (!result.Ok)
  73. {
  74. ok = false;
  75. break;
  76. }
  77. await Task.Delay(300, ct); // quota 보호 (청크 간)
  78. }
  79. }
  80. else
  81. {
  82. // 증분 수집: 신규분이 최신 페이지에 몰려 있으므로 전체 기존 페이지를 만나면 조기 종료
  83. var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, beginDate, todayKst, allowEarlyStop: true, ct);
  84. inserted = result.Inserted;
  85. skipped = result.Skipped;
  86. ok = result.Ok;
  87. }
  88. Logger.LogInformation("[{Job}] 완료 — 범위=[{Begin}~{Today}], inserted={Inserted}, skipped={Skipped}, 초기백필={Initial}",
  89. JobName, beginDate, todayKst, inserted, skipped, isInitialBackfill);
  90. return ok;
  91. }
  92. /// <summary>
  93. /// [bgnDe, endDe] 접수일 범위를 desc 정렬로 페이징 수집. 신규만 삽입(rcept_no idempotent).
  94. /// allowEarlyStop=true 면 "페이지 전체 기존 + 1페이지 아님" 조건에서 조기 종료 (DartBackfill.ShouldStopEarly).
  95. /// API 오류 status 는 Ok=false 로 반환해 베이스 재시도로 연결한다.
  96. /// </summary>
  97. private async Task<(bool Ok, int Inserted, int Skipped)> SyncRangeAsync(
  98. IAppDbContext db,
  99. HttpClient client,
  100. string baseUrl,
  101. string apiKey,
  102. DateOnly bgnDe,
  103. DateOnly endDe,
  104. bool allowEarlyStop,
  105. CancellationToken ct)
  106. {
  107. var inserted = 0;
  108. var skipped = 0;
  109. var page = 1;
  110. var totalPage = 1;
  111. while (page <= totalPage && page <= MaxPages)
  112. {
  113. ct.ThrowIfCancellationRequested();
  114. var url = $"{baseUrl}/api/list.json?crtfc_key={apiKey}&bgn_de={bgnDe:yyyyMMdd}&end_de={endDe:yyyyMMdd}&page_no={page}&page_count={PageCount}";
  115. var json = await OpenDartHttp.GetStringWithRetryAsync(client, url, Logger, ct);
  116. var result = DartDisclosureParser.Parse(json);
  117. if (result.Status == "013")
  118. {
  119. Logger.LogInformation("[{Job}] 조회된 데이터 없음(013) — 종료. 범위=[{Begin}~{End}]", JobName, bgnDe, endDe);
  120. break;
  121. }
  122. if (result.Status != "000")
  123. {
  124. Logger.LogError("[{Job}] OpenDART 오류 status={Status}({Label}) message={Message} (page={Page}, 범위=[{Begin}~{End}])",
  125. JobName, result.Status, DartDisclosureParser.DescribeStatus(result.Status), result.Message, page, bgnDe, endDe);
  126. return (false, inserted, skipped);
  127. }
  128. totalPage = result.TotalPage > 0 ? result.TotalPage : 1;
  129. if (result.Rows.Count == 0)
  130. {
  131. break;
  132. }
  133. var rceptNos = result.Rows.Select(c => c.RceptNo).ToList();
  134. var existing = (await db.Disclosure.AsNoTracking()
  135. .Where(c => rceptNos.Contains(c.RceptNo))
  136. .Select(c => c.RceptNo)
  137. .ToListAsync(ct)).ToHashSet();
  138. var newCount = 0;
  139. foreach (var row in result.Rows)
  140. {
  141. if (existing.Contains(row.RceptNo))
  142. {
  143. skipped++;
  144. continue;
  145. }
  146. var created = Disclosure.Create(row.RceptNo, row.CorpCode, row.CorpName, row.StockCode, row.CorpCls, row.ReportNm, row.FlrNm, row.RceptDt, row.Rm);
  147. await db.Disclosure.AddAsync(created, ct);
  148. existing.Add(row.RceptNo); // 같은 실행 내 중복 접수번호 방어
  149. inserted++;
  150. newCount++;
  151. }
  152. await db.SaveChangesAsync(ct);
  153. Logger.LogInformation("[{Job}] page {Page}/{TotalPage} 처리 — rows={Rows}, 신규={New} (누적 inserted={Inserted}, skipped={Skipped})",
  154. JobName, page, totalPage, result.Rows.Count, newCount, inserted, skipped);
  155. if (allowEarlyStop && DartBackfill.ShouldStopEarly(page, result.Rows.Count, newCount))
  156. {
  157. Logger.LogInformation("[{Job}] 조기 종료 — page {Page} 전체가 기수집 건 (이후 페이지는 더 오래된 기수집분)", JobName, page);
  158. return (true, inserted, skipped);
  159. }
  160. page++;
  161. if (page <= totalPage)
  162. {
  163. await Task.Delay(300, ct); // quota 보호
  164. }
  165. }
  166. if (totalPage > MaxPages)
  167. {
  168. Logger.LogWarning("[{Job}] 페이지 상한 도달 — totalPage={TotalPage}, 상한={MaxPages}, 미수집 잔여 {Remaining}페이지 (정렬 desc 라 가장 오래된 공시부터 유실). 범위=[{Begin}~{End}] 를 좁혀 재수집 필요",
  169. JobName, totalPage, MaxPages, totalPage - MaxPages, bgnDe, endDe);
  170. }
  171. return (true, inserted, skipped);
  172. }
  173. }