DisclosureSyncService.cs 4.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123
  1. using Application.Abstractions.Data;
  2. using Domain.Entities.Stocks;
  3. using Microsoft.EntityFrameworkCore;
  4. using Microsoft.Extensions.DependencyInjection;
  5. using Microsoft.Extensions.Logging;
  6. using Microsoft.Extensions.Options;
  7. using SharedKernel;
  8. namespace Infrastructure.StockData;
  9. /// <summary>
  10. /// 전자공시(DART) 공시 목록 수집 (OpenDART list.json) — 기본 08:30 KST 실행. 접수일 [today-BackfillDays, today] 범위를
  11. /// page_count=100 으로 페이징(total_page 만큼 순회)하며 수집한다. KRX 배치처럼 날짜별 반복이 아니라 DART 는 날짜범위 페이징이므로
  12. /// KrxBackfill 대신 자체 페이징 루프를 쓴다. 접수번호(rcept_no) 로 upsert 하되 이미 있으면 skip(idempotent) — 갱신 없이 신규만 삽입.
  13. /// status "013"(데이터 없음) 은 정상 종료, "000" 외 다른 status 는 오류로 보고 후 중단. ApiKey 미설정 시 로그만 남기고 skip.
  14. /// list.json 이 stock_code 를 포함하므로 corpCode.xml 매핑 불필요.
  15. /// </summary>
  16. internal sealed class DisclosureSyncService(
  17. IServiceScopeFactory scopeFactory,
  18. IHttpClientFactory httpClientFactory,
  19. IOptions<AppSettings> settings,
  20. ILogger<DisclosureSyncService> logger
  21. ) : DailyScheduledService(logger)
  22. {
  23. private const int PageCount = 100;
  24. private const int MaxPages = 500; // 안전 상한 (total_page 무한루프 방지)
  25. protected override string JobName => "DisclosureSync";
  26. protected override TimeOnly TargetTime => ParseTime(settings.Value.OpenDart.DisclosureSyncTime, new TimeOnly(8, 30));
  27. protected override int MaxRetryCount => 2;
  28. protected override TimeSpan RetryDelay => TimeSpan.FromHours(1);
  29. protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
  30. {
  31. var cfg = settings.Value.OpenDart;
  32. if (string.IsNullOrWhiteSpace(cfg.ApiKey))
  33. {
  34. Logger.LogWarning("[{Job}] OpenDart:ApiKey 미설정 — 수집 skip", JobName);
  35. return true;
  36. }
  37. using var scope = scopeFactory.CreateScope();
  38. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  39. var client = httpClientFactory.CreateClient(OpenDartHttp.ClientName);
  40. var days = cfg.BackfillDays > 0 ? cfg.BackfillDays : 30;
  41. var beginDate = todayKst.AddDays(-days);
  42. var baseUrl = cfg.BaseUrl.TrimEnd('/');
  43. var inserted = 0;
  44. var skipped = 0;
  45. var page = 1;
  46. var totalPage = 1;
  47. while (page <= totalPage && page <= MaxPages)
  48. {
  49. ct.ThrowIfCancellationRequested();
  50. var url = $"{baseUrl}/api/list.json?crtfc_key={cfg.ApiKey}&bgn_de={beginDate:yyyyMMdd}&end_de={todayKst:yyyyMMdd}&page_no={page}&page_count={PageCount}";
  51. var json = await OpenDartHttp.GetStringWithRetryAsync(client, url, Logger, ct);
  52. var result = DartDisclosureParser.Parse(json);
  53. if (result.Status == "013")
  54. {
  55. Logger.LogInformation("[{Job}] 조회된 데이터 없음(013) — 종료. 범위=[{Begin}~{Today}]", JobName, beginDate, todayKst);
  56. break;
  57. }
  58. if (result.Status != "000")
  59. {
  60. Logger.LogError("[{Job}] OpenDART 오류 status={Status} message={Message} (page={Page})", JobName, result.Status, result.Message, page);
  61. return false;
  62. }
  63. totalPage = result.TotalPage > 0 ? result.TotalPage : 1;
  64. if (result.Rows.Count == 0)
  65. {
  66. break;
  67. }
  68. var rceptNos = result.Rows.Select(c => c.RceptNo).ToList();
  69. var existing = (await db.Disclosure.AsNoTracking()
  70. .Where(c => rceptNos.Contains(c.RceptNo))
  71. .Select(c => c.RceptNo)
  72. .ToListAsync(ct)).ToHashSet();
  73. foreach (var row in result.Rows)
  74. {
  75. if (existing.Contains(row.RceptNo))
  76. {
  77. skipped++;
  78. continue;
  79. }
  80. var created = Disclosure.Create(row.RceptNo, row.CorpCode, row.CorpName, row.StockCode, row.CorpCls, row.ReportNm, row.FlrNm, row.RceptDt, row.Rm);
  81. await db.Disclosure.AddAsync(created, ct);
  82. existing.Add(row.RceptNo); // 같은 실행 내 중복 접수번호 방어
  83. inserted++;
  84. }
  85. await db.SaveChangesAsync(ct);
  86. Logger.LogInformation("[{Job}] page {Page}/{TotalPage} 처리 — rows={Rows} (누적 inserted={Inserted}, skipped={Skipped})",
  87. JobName, page, totalPage, result.Rows.Count, inserted, skipped);
  88. page++;
  89. if (page <= totalPage)
  90. {
  91. await Task.Delay(300, ct); // quota 보호
  92. }
  93. }
  94. Logger.LogInformation("[{Job}] 완료 — 범위=[{Begin}~{Today}], inserted={Inserted}, skipped={Skipped}",
  95. JobName, beginDate, todayKst, inserted, skipped);
  96. return true;
  97. }
  98. }