| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210 |
- using Application.Abstractions.Data;
- using Application.Helpers;
- using Domain.Entities.Stocks;
- using Microsoft.EntityFrameworkCore;
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Logging;
- using Microsoft.Extensions.Options;
- using SharedKernel;
- namespace Infrastructure.StockData;
- /// <summary>
- /// 전자공시(DART) 공시 목록 수집 (OpenDART list.json) — 기본 08:30 KST 실행. 접수일 [today-BackfillDays, today] 범위를
- /// page_count=100 으로 페이징(total_page 만큼 순회)하며 수집한다. KRX 배치처럼 날짜별 반복이 아니라 DART 는 날짜범위 페이징이므로
- /// KrxBackfill 대신 자체 페이징 루프를 쓴다. 접수번호(rcept_no) 로 upsert 하되 이미 있으면 skip(idempotent) — 갱신 없이 신규만 삽입.
- /// • 스펙 클램프: corp_code 미지정 검색은 최근 3개월만 제공 — BackfillDays 가 90 초과여도 창을 3개월로 클램프(경고).
- /// • 초기 백필(창의 가장 오래된 주에 적재분 0건): bgn_de/end_de 를 주 단위로 분할 순회해 5만 행(MaxPages×100) 상한을 회피.
- /// • 증분 수집: 한 페이지 전체가 기존 건이면(1페이지 제외) 그보다 오래된 페이지도 기수집으로 보고 조기 종료 — 매일 전량 재페이징 방지.
- /// • 페이지 상한 도달 시 잔여 페이지 수를 경고 (정렬 desc 라 오래된 공시부터 유실).
- /// status "013"(데이터 없음) 은 정상 종료, "000" 외 다른 status 는 오류로 보고 후 중단(재시도 대상 — DescribeStatus 라벨 로깅).
- /// ApiKey 미설정 시 로그만 남기고 skip. list.json 이 stock_code 를 포함하므로 corpCode.xml 매핑 불필요.
- /// </summary>
- internal sealed class DisclosureSyncService(
- IServiceScopeFactory scopeFactory,
- IHttpClientFactory httpClientFactory,
- IOptions<AppSettings> settings,
- ILogger<DisclosureSyncService> logger
- ) : DailyScheduledService(logger)
- {
- private const int PageCount = 100;
- private const int MaxPages = 500; // 안전 상한 (total_page 무한루프 방지)
- protected override string JobName => "DisclosureSync";
- protected override TimeOnly TargetTime => ParseTime(settings.Value.OpenDart.DisclosureSyncTime, new TimeOnly(8, 30));
- protected override int MaxRetryCount => 2;
- protected override TimeSpan RetryDelay => TimeSpan.FromHours(1);
- protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
- {
- var cfg = settings.Value.OpenDart;
- if (string.IsNullOrWhiteSpace(cfg.ApiKey))
- {
- Logger.LogWarning("[{Job}] OpenDart:ApiKey 미설정 — 수집 skip", JobName);
- return true;
- }
- using var scope = scopeFactory.CreateScope();
- var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
- var client = httpClientFactory.CreateClient(OpenDartHttp.ClientName);
- var days = cfg.BackfillDays > 0 ? cfg.BackfillDays : 30;
- var beginDate = todayKst.AddDays(-days);
- // corp_code 미지정 list.json 은 스펙상 최근 3개월만 제공 — 초과 설정 시 클램프
- var (clampedBegin, clamped) = DartBackfill.ClampToSpecWindow(beginDate, todayKst);
- if (clamped)
- {
- Logger.LogWarning("[{Job}] BackfillDays={Days} 가 스펙 한도(최근 {Months}개월)를 초과 — 창 시작일을 {Begin} 으로 클램프",
- JobName, days, DartBackfill.SpecWindowMonths, clampedBegin);
- beginDate = clampedBegin;
- }
- var baseUrl = cfg.BaseUrl.TrimEnd('/');
- // 초기 백필 판정 — 창의 가장 오래된 주(7일)에 적재된 공시가 하나도 없으면 미커버로 간주 (영업일마다 수백 건이 접수되므로 안전)
- var oldestWeekEnd = beginDate.AddDays(7);
- var isInitialBackfill = !await db.Disclosure.AsNoTracking().AnyAsync(c => c.RceptDt >= beginDate && c.RceptDt < oldestWeekEnd, ct);
- var inserted = 0;
- var skipped = 0;
- var ok = true;
- if (isInitialBackfill)
- {
- // 최초 적재: 창 전체가 5만 행(MaxPages×PageCount) 상한을 넘을 수 있어 주 단위로 분할 순회 (사업보고서 시즌 30일 창이 5만 초과 현실적)
- var chunks = DartBackfill.SplitWeekly(beginDate, todayKst);
- Logger.LogInformation("[{Job}] 초기 백필 — 범위=[{Begin}~{Today}] 를 주 단위 {Chunks}개 창으로 분할 수집", JobName, beginDate, todayKst, chunks.Count);
- foreach (var (chunkBegin, chunkEnd) in chunks)
- {
- var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, chunkBegin, chunkEnd, allowEarlyStop: false, ct);
- inserted += result.Inserted;
- skipped += result.Skipped;
- if (!result.Ok)
- {
- ok = false;
- break;
- }
- await Task.Delay(300, ct); // quota 보호 (청크 간)
- }
- }
- else
- {
- // 증분 수집: 신규분이 최신 페이지에 몰려 있으므로 전체 기존 페이지를 만나면 조기 종료
- var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, beginDate, todayKst, allowEarlyStop: true, ct);
- inserted = result.Inserted;
- skipped = result.Skipped;
- ok = result.Ok;
- }
- Logger.LogInformation("[{Job}] 완료 — 범위=[{Begin}~{Today}], inserted={Inserted}, skipped={Skipped}, 초기백필={Initial}",
- JobName, beginDate, todayKst, inserted, skipped, isInitialBackfill);
- return ok;
- }
- /// <summary>
- /// [bgnDe, endDe] 접수일 범위를 desc 정렬로 페이징 수집. 신규만 삽입(rcept_no idempotent).
- /// allowEarlyStop=true 면 "페이지 전체 기존 + 1페이지 아님" 조건에서 조기 종료 (DartBackfill.ShouldStopEarly).
- /// API 오류 status 는 Ok=false 로 반환해 베이스 재시도로 연결한다.
- /// </summary>
- private async Task<(bool Ok, int Inserted, int Skipped)> SyncRangeAsync(
- IAppDbContext db,
- HttpClient client,
- string baseUrl,
- string apiKey,
- DateOnly bgnDe,
- DateOnly endDe,
- bool allowEarlyStop,
- CancellationToken ct)
- {
- var inserted = 0;
- var skipped = 0;
- var page = 1;
- var totalPage = 1;
- while (page <= totalPage && page <= MaxPages)
- {
- ct.ThrowIfCancellationRequested();
- var url = $"{baseUrl}/api/list.json?crtfc_key={apiKey}&bgn_de={bgnDe:yyyyMMdd}&end_de={endDe:yyyyMMdd}&page_no={page}&page_count={PageCount}";
- var json = await OpenDartHttp.GetStringWithRetryAsync(client, url, Logger, ct);
- var result = DartDisclosureParser.Parse(json);
- if (result.Status == "013")
- {
- Logger.LogInformation("[{Job}] 조회된 데이터 없음(013) — 종료. 범위=[{Begin}~{End}]", JobName, bgnDe, endDe);
- break;
- }
- if (result.Status != "000")
- {
- Logger.LogError("[{Job}] OpenDART 오류 status={Status}({Label}) message={Message} (page={Page}, 범위=[{Begin}~{End}])",
- JobName, result.Status, DartDisclosureParser.DescribeStatus(result.Status), result.Message, page, bgnDe, endDe);
- return (false, inserted, skipped);
- }
- totalPage = result.TotalPage > 0 ? result.TotalPage : 1;
- if (result.Rows.Count == 0)
- {
- break;
- }
- var rceptNos = result.Rows.Select(c => c.RceptNo).ToList();
- var existing = (await db.Disclosure.AsNoTracking()
- .Where(c => rceptNos.Contains(c.RceptNo))
- .Select(c => c.RceptNo)
- .ToListAsync(ct)).ToHashSet();
- var newCount = 0;
- foreach (var row in result.Rows)
- {
- if (existing.Contains(row.RceptNo))
- {
- skipped++;
- continue;
- }
- var created = Disclosure.Create(row.RceptNo, row.CorpCode, row.CorpName, row.StockCode, row.CorpCls, row.ReportNm, row.FlrNm, row.RceptDt, row.Rm);
- await db.Disclosure.AddAsync(created, ct);
- existing.Add(row.RceptNo); // 같은 실행 내 중복 접수번호 방어
- inserted++;
- newCount++;
- }
- await db.SaveChangesAsync(ct);
- Logger.LogInformation("[{Job}] page {Page}/{TotalPage} 처리 — rows={Rows}, 신규={New} (누적 inserted={Inserted}, skipped={Skipped})",
- JobName, page, totalPage, result.Rows.Count, newCount, inserted, skipped);
- if (allowEarlyStop && DartBackfill.ShouldStopEarly(page, result.Rows.Count, newCount))
- {
- Logger.LogInformation("[{Job}] 조기 종료 — page {Page} 전체가 기수집 건 (이후 페이지는 더 오래된 기수집분)", JobName, page);
- return (true, inserted, skipped);
- }
- page++;
- if (page <= totalPage)
- {
- await Task.Delay(300, ct); // quota 보호
- }
- }
- if (totalPage > MaxPages)
- {
- Logger.LogWarning("[{Job}] 페이지 상한 도달 — totalPage={TotalPage}, 상한={MaxPages}, 미수집 잔여 {Remaining}페이지 (정렬 desc 라 가장 오래된 공시부터 유실). 범위=[{Begin}~{End}] 를 좁혀 재수집 필요",
- JobName, totalPage, MaxPages, totalPage - MaxPages, bgnDe, endDe);
- }
- return (true, inserted, skipped);
- }
- }
|