SeibroDividendSyncService.cs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221
  1. using Application.Abstractions.Data;
  2. using Application.Helpers;
  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. /// SEIBro 배당·권리 수집 (Wave 1 ★P0) — 기본 06:40 KST(IssuerSync 06:30 다음 슬롯), `Seibro:DividendSync` 게이트.
  11. /// 세 단계가 하나의 기업(Corp) 카테고리 예산(SeibroQuota, Seibro:CorpBudget)을 공유한다.
  12. ///
  13. /// a) DividendSchedule 날짜 스윕: getDivSchedulInfo(BEGIN_STD_DT=day, EXPRY 생략 → 그날만, custno 미지정 → 전체 회사).
  14. /// 3년 창을 최신일→과거로 훑으며 미적재일만 수집한다.
  15. /// ⚠️ 날짜 스윕에 KrxBackfill(영업일·휴장 스킵)을 쓰지 않는다 — 배당 권리기준일 상당수가 분기말/연말(0331·0630·0930·1231)이며
  16. /// 이 날짜가 토·일(예: 20161231 토, 20180331 토, 20180630 토)인 경우가 실제 샘플에 존재한다. 영업일만 훑으면 이런 기준일을 통째로 놓친다.
  17. /// → 주말·휴장 포함 전 캘린더일을 훑는다(3년 ≈ 1,095콜, Corp 예산 40,000 내 충분).
  18. /// b) Dividend rolling: 대상 = DividendSchedule 에 등장한 DISTINCT IssucoCustno 중 Dividend 미수집(우선)·stale 순.
  19. /// getDivInfo(ISSUCO_CUSTNO, BEGIN_STD_DT=3년전, EXPRY_STD_DT=오늘) 1콜 → 여러 배당 반환 → Dividend upsert(Isin+RgtStdDt).
  20. /// getDivInfo 응답엔 ISSUCO_CUSTNO 가 없으므로(샘플 확인) Isin 으로만 upsert 한다.
  21. /// c) RightsBaseDate rolling: 대상 = Stock.IssucoCustno non-null 중 RightsBaseDate 미수집(우선)·stale 순.
  22. /// getStddtInfo(ISSUCO_CUSTNO, BEGIN_STD_DT=3년전, EXPRY_STD_DT=오늘) → RightsBaseDate upsert.
  23. ///
  24. /// 실패/0행 = 정상 빈결과(다음 실행 재시도). HTTP·파싱 오류 = 경고 + false(RetryDelay 후 재시도).
  25. /// quota 소진 = 정상 종료(true) — 일일 예산은 KST 자정 롤오버로만 회복. ApiKey 미설정 시 로그만 남기고 skip.
  26. /// </summary>
  27. internal sealed class SeibroDividendSyncService(
  28. IServiceScopeFactory scopeFactory,
  29. IHttpClientFactory httpClientFactory,
  30. SeibroQuota quota,
  31. IOptions<AppSettings> settings,
  32. ILogger<SeibroDividendSyncService> logger
  33. ) : DailyScheduledService(logger)
  34. {
  35. protected override string JobName => "SeibroDividendSync";
  36. protected override TimeOnly TargetTime => ParseTime(settings.Value.Seibro.DividendSyncTime, new TimeOnly(6, 40));
  37. protected override int MaxRetryCount => 2;
  38. protected override TimeSpan RetryDelay => TimeSpan.FromHours(2);
  39. protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
  40. {
  41. var cfg = settings.Value.Seibro;
  42. if (string.IsNullOrWhiteSpace(cfg.ApiKey))
  43. {
  44. Logger.LogWarning("[{Job}] Seibro:ApiKey 미설정 — 수집 skip", JobName);
  45. return true;
  46. }
  47. using var scope = scopeFactory.CreateScope();
  48. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  49. var client = httpClientFactory.CreateClient(SeibroHttp.ClientName);
  50. var years = cfg.BackfillYears > 0 ? cfg.BackfillYears : 3;
  51. var startDate = todayKst.AddYears(-years);
  52. try
  53. {
  54. if (!await SweepSchedulesAsync(db, client, cfg, startDate, todayKst, ct))
  55. {
  56. return true; // quota 소진 — 오늘은 더 진행 불가 (정상 종료)
  57. }
  58. if (!await RollDividendsAsync(db, client, cfg, startDate, todayKst, ct))
  59. {
  60. return true;
  61. }
  62. await RollRightsBaseDatesAsync(db, client, cfg, startDate, todayKst, ct);
  63. return true;
  64. }
  65. catch (Exception ex) when (ex is HttpRequestException or System.Xml.XmlException or FormatException)
  66. {
  67. Logger.LogWarning(ex, "[{Job}] SEIBro 호출/파싱 실패 — run 중단, {Delay} 후 재시도 (최대 {Max}회)", JobName, RetryDelay, MaxRetryCount);
  68. return false;
  69. }
  70. }
  71. /// <summary>a) getDivSchedulInfo 날짜 스윕 (전 캘린더일, 최신일→과거) — 미적재일만 수집. quota 소진 시 false.</summary>
  72. private async Task<bool> SweepSchedulesAsync(IAppDbContext db, HttpClient client, AppSettings.SeibroSection cfg, DateOnly startDate, DateOnly today, CancellationToken ct)
  73. {
  74. var fetched = 0;
  75. var first = true;
  76. for (var day = today; day >= startDate; day = day.AddDays(-1))
  77. {
  78. ct.ThrowIfCancellationRequested();
  79. if (await db.DividendSchedule.AsNoTracking().AnyAsync(c => c.RgtStdDt == day, ct))
  80. {
  81. continue; // 이미 적재된 날 — resumable
  82. }
  83. if (!quota.TryConsume(SeibroCategory.Corp, 1, cfg.CorpBudget))
  84. {
  85. Logger.LogWarning("[{Job}] 기업 카테고리 일일 예산({Budget}) 소진 — getDivSchedulInfo {Day} 부터 중단", JobName, cfg.CorpBudget, day);
  86. return false;
  87. }
  88. if (!first && cfg.DelayMs > 0)
  89. {
  90. await Task.Delay(cfg.DelayMs, ct);
  91. }
  92. first = false;
  93. var xml = await SeibroHttp.GetStringWithRetryAsync(client, cfg.BaseUrl, cfg.ApiKey, "getDivSchedulInfo", [new("BEGIN_STD_DT", day.ToString("yyyyMMdd"))], Logger, ct);
  94. var rows = SeibroDivSchedulParser.Parse(SeibroXml.Parse(xml));
  95. var (inserted, updated) = await SeibroDividendImport.UpsertSchedulesAsync(db, rows, ct);
  96. Logger.LogInformation("[{Job}] getDivSchedulInfo {Day} rows={Rows} — inserted={Inserted}, updated={Updated}", JobName, day, rows.Count, inserted, updated);
  97. fetched++;
  98. }
  99. Logger.LogInformation("[{Job}] 배당일정 스윕 완료 — 창=[{Start}~{End}], fetch={Fetched}일", JobName, startDate, today, fetched);
  100. return true;
  101. }
  102. /// <summary>b) getDivInfo rolling — DividendSchedule 에 등장한 회사번호 중 미수집·stale 우선. quota 잔여분 내. 예산 소진 시 false.</summary>
  103. private async Task<bool> RollDividendsAsync(IAppDbContext db, HttpClient client, AppSettings.SeibroSection cfg, DateOnly startDate, DateOnly today, CancellationToken ct)
  104. {
  105. var scheduleCustnos = await db.DividendSchedule.AsNoTracking().Select(c => c.IssucoCustno).Distinct().ToListAsync(ct);
  106. if (scheduleCustnos.Count == 0)
  107. {
  108. Logger.LogInformation("[{Job}] getDivInfo 대상 없음 (DividendSchedule 비어 있음)", JobName);
  109. return true;
  110. }
  111. // 회사번호 → 해당 회사 ISIN 들의 Dividend 최신 UpdatedAt (미수집이면 null → 최우선)
  112. var custnoIsins = await db.Stock.AsNoTracking().Where(c => c.IssucoCustno != null && c.ISIN != null)
  113. .Select(c => new { Custno = c.IssucoCustno!.Value, c.ISIN })
  114. .ToListAsync(ct);
  115. var isinsByCustno = custnoIsins.GroupBy(c => c.Custno).ToDictionary(g => g.Key, g => g.Select(c => c.ISIN!).ToList());
  116. var dividendUpdatedByIsin = await db.Dividend.AsNoTracking().GroupBy(c => c.Isin).Select(g => new { Isin = g.Key, Last = g.Max(c => c.UpdatedAt) }).ToDictionaryAsync(c => c.Isin, c => c.Last, ct);
  117. var targets = scheduleCustnos.Select(custno => {
  118. DateTime? last = null;
  119. if (isinsByCustno.TryGetValue(custno, out var isins))
  120. {
  121. var times = isins.Where(dividendUpdatedByIsin.ContainsKey).Select(i => dividendUpdatedByIsin[i]).ToList();
  122. // 회사의 모든 ISIN 이 최소 1회 수집됐을 때만 stale(가장 오래된 시각) 로 본다 — 하나라도 미수집이면 미수집(null=최우선)
  123. last = isins.Count > 0 && times.Count == isins.Count ? times.Min() : null;
  124. }
  125. return (Key: custno, LastUpdatedAt: last);
  126. }).ToList();
  127. var maxPerRun = quota.Remaining(SeibroCategory.Corp, cfg.CorpBudget);
  128. var beginStr = startDate.ToString("yyyyMMdd");
  129. var expiryStr = today.ToString("yyyyMMdd");
  130. var quotaExhausted = false;
  131. var processed = await SeibroRollingSweep.RunAsync(
  132. targets: targets,
  133. fetchAndUpsert: async (custno, token) => {
  134. if (!quota.TryConsume(SeibroCategory.Corp, 1, cfg.CorpBudget))
  135. {
  136. quotaExhausted = true;
  137. return;
  138. }
  139. var xml = await SeibroHttp.GetStringWithRetryAsync(client, cfg.BaseUrl, cfg.ApiKey, "getDivInfo", [new("ISSUCO_CUSTNO", custno.ToString()), new("BEGIN_STD_DT", beginStr), new("EXPRY_STD_DT", expiryStr)], Logger, token);
  140. var rows = SeibroDivInfoParser.Parse(SeibroXml.Parse(xml));
  141. var (inserted, updated) = await SeibroDividendImport.UpsertDividendsAsync(db, rows, token);
  142. Logger.LogInformation("[{Job}] getDivInfo ISSUCO_CUSTNO={Custno} rows={Rows} — inserted={Inserted}, updated={Updated}", JobName, custno, rows.Count, inserted, updated);
  143. },
  144. maxPerRun: maxPerRun,
  145. delayMs: cfg.DelayMs,
  146. ct: ct);
  147. Logger.LogInformation("[{Job}] getDivInfo rolling 완료 — 대상={Targets}, 호출={Processed} (maxPerRun={Max})", JobName, targets.Count, processed, maxPerRun);
  148. return !quotaExhausted;
  149. }
  150. /// <summary>c) getStddtInfo rolling — Stock.IssucoCustno non-null 중 미수집·stale 우선. quota 잔여분 내.</summary>
  151. private async Task RollRightsBaseDatesAsync(IAppDbContext db, HttpClient client, AppSettings.SeibroSection cfg, DateOnly startDate, DateOnly today, CancellationToken ct)
  152. {
  153. var custnos = await db.Stock.AsNoTracking().Where(c => c.IssucoCustno != null).Select(c => c.IssucoCustno!.Value).Distinct().ToListAsync(ct);
  154. if (custnos.Count == 0)
  155. {
  156. Logger.LogInformation("[{Job}] getStddtInfo 대상 없음 (Stock.IssucoCustno 미스탬핑)", JobName);
  157. return;
  158. }
  159. var lastByCustno = await db.RightsBaseDate.AsNoTracking().GroupBy(c => c.IssucoCustno).Select(g => new { Custno = g.Key, Last = g.Max(c => c.UpdatedAt) }).ToDictionaryAsync(c => c.Custno, c => c.Last, ct);
  160. var targets = custnos.Select(custno => (Key: custno, LastUpdatedAt: lastByCustno.TryGetValue(custno, out var last) ? (DateTime?)last : null)).ToList();
  161. var maxPerRun = quota.Remaining(SeibroCategory.Corp, cfg.CorpBudget);
  162. var beginStr = startDate.ToString("yyyyMMdd");
  163. var expiryStr = today.ToString("yyyyMMdd");
  164. var processed = await SeibroRollingSweep.RunAsync(
  165. targets: targets,
  166. fetchAndUpsert: async (custno, token) => {
  167. if (!quota.TryConsume(SeibroCategory.Corp, 1, cfg.CorpBudget))
  168. {
  169. return;
  170. }
  171. var xml = await SeibroHttp.GetStringWithRetryAsync(client, cfg.BaseUrl, cfg.ApiKey, "getStddtInfo", [new("ISSUCO_CUSTNO", custno.ToString()), new("BEGIN_STD_DT", beginStr), new("EXPRY_STD_DT", expiryStr)], Logger, token);
  172. var rows = SeibroStddtParser.Parse(SeibroXml.Parse(xml));
  173. var (inserted, updated) = await SeibroDividendImport.UpsertRightsBaseDatesAsync(db, rows, token);
  174. Logger.LogInformation("[{Job}] getStddtInfo ISSUCO_CUSTNO={Custno} rows={Rows} — inserted={Inserted}, updated={Updated}", JobName, custno, rows.Count, inserted, updated);
  175. },
  176. maxPerRun: maxPerRun,
  177. delayMs: cfg.DelayMs,
  178. ct: ct);
  179. Logger.LogInformation("[{Job}] getStddtInfo rolling 완료 — 대상={Targets}, 호출={Processed} (maxPerRun={Max})", JobName, targets.Count, processed, maxPerRun);
  180. }
  181. }