MacroDataSyncService.cs 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204
  1. using Application.Abstractions.Data;
  2. using Application.Helpers;
  3. using Domain.Entities.Stocks;
  4. using Domain.Entities.Stocks.ValueObject;
  5. using Microsoft.EntityFrameworkCore;
  6. using Microsoft.Extensions.DependencyInjection;
  7. using Microsoft.Extensions.Logging;
  8. using Microsoft.Extensions.Options;
  9. using SharedKernel;
  10. namespace Infrastructure.StockData;
  11. /// <summary>
  12. /// 거시·환율 수집 (한국수출입은행 OpenAPI) — AP01 현재환율 + AP02 대출금리 + AP03 국제금리. 기본 11:30 KST 실행
  13. /// (koreaexim 은 영업일 ~11시 이후 반영). 공용 KrxBackfill 로 [today-MacroBackfillDays, 직전 영업일] 을 최신일부터 과거로 훑으며
  14. /// 미적재일만 채운다 — quota 보호를 위해 1회 실행당 BackfillMaxPerRun(기본 60)일까지만 fetch. 각 날짜는 세 data 종류를
  15. /// (키가 설정된 것만) 모두 수집해 ExchangeRate / InterestRate upsert (UQ = (CurrencyUnit,TradeDate) / (RateType,ItemName,TradeDate)).
  16. /// 빈 배열([]) — 주말/공휴일·오전 미반영 — 은 "데이터 없음" 으로 취급(그 날은 미적재 유지, 다음 실행 때 재시도). 키 전부 미설정 시 skip.
  17. /// existsForDate 판정은 세 종류 중 하나라도 그 날짜에 적재됐는지로 본다 (부분 성공 후 재실행 시 남은 종류를 다시 채울 수 있게 upsert 로 안전).
  18. /// </summary>
  19. internal sealed class MacroDataSyncService(
  20. IServiceScopeFactory scopeFactory,
  21. IHttpClientFactory httpClientFactory,
  22. IOptions<AppSettings> settings,
  23. ILogger<MacroDataSyncService> logger
  24. ) : DailyScheduledService(logger)
  25. {
  26. protected override string JobName => "MacroDataSync";
  27. protected override TimeOnly TargetTime => ParseTime(settings.Value.KoreaExim.MacroSyncTime, new TimeOnly(11, 30));
  28. protected override int MaxRetryCount => 2;
  29. protected override TimeSpan RetryDelay => TimeSpan.FromHours(1);
  30. protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
  31. {
  32. var cfg = settings.Value.KoreaExim;
  33. var hasExchange = !string.IsNullOrWhiteSpace(cfg.ExchangeKey);
  34. var hasLoan = !string.IsNullOrWhiteSpace(cfg.LoanRateKey);
  35. var hasIntl = !string.IsNullOrWhiteSpace(cfg.IntlRateKey);
  36. if (!hasExchange && !hasLoan && !hasIntl)
  37. {
  38. Logger.LogWarning("[{Job}] KoreaExim 키(ExchangeKey/LoanRateKey/IntlRateKey) 전부 미설정 — 수집 skip", JobName);
  39. return true;
  40. }
  41. using var scope = scopeFactory.CreateScope();
  42. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  43. var client = httpClientFactory.CreateClient(KoreaEximHttp.ClientName);
  44. var endDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, todayKst.AddDays(1), ct);
  45. var days = cfg.MacroBackfillDays > 0 ? cfg.MacroBackfillDays : 30;
  46. var startDate = todayKst.AddDays(-days);
  47. var holidays = (await db.MarketHoliday.AsNoTracking()
  48. .Where(c => c.Date >= startDate && c.Date <= endDate)
  49. .Select(c => c.Date)
  50. .ToListAsync(ct)).ToHashSet();
  51. var maxPerRun = cfg.BackfillMaxPerRun > 0 ? cfg.BackfillMaxPerRun : 60;
  52. var baseUrl = cfg.BaseUrl.TrimEnd('/');
  53. var fetched = await KrxBackfill.RunAsync(
  54. existsForDate: (day, token) => ExistsForDateAsync(db, day, hasExchange, token),
  55. fetchAndUpsertForDate: (day, token) => FetchAndUpsertAsync(db, client, baseUrl, cfg, hasExchange, hasLoan, hasIntl, day, token),
  56. startDate: startDate,
  57. endDate: endDate,
  58. holidays: holidays,
  59. maxPerRun: maxPerRun,
  60. delayMs: 300,
  61. ct: ct);
  62. Logger.LogInformation("[{Job}] 완료 — 창=[{Start}~{End}], 이번 실행 fetch={Fetched}일 (maxPerRun={Max})",
  63. JobName, startDate, endDate, fetched, maxPerRun);
  64. return true;
  65. }
  66. /// <summary>해당 날짜 적재 여부 — 환율 키가 있으면 ExchangeRate, 없으면 InterestRate 로 판정</summary>
  67. private static async Task<bool> ExistsForDateAsync(IAppDbContext db, DateOnly day, bool hasExchange, CancellationToken ct)
  68. {
  69. if (hasExchange)
  70. {
  71. return await db.ExchangeRate.AsNoTracking().AnyAsync(c => c.TradeDate == day, ct);
  72. }
  73. return await db.InterestRate.AsNoTracking().AnyAsync(c => c.TradeDate == day, ct);
  74. }
  75. private async Task FetchAndUpsertAsync(
  76. IAppDbContext db,
  77. HttpClient client,
  78. string baseUrl,
  79. AppSettings.KoreaEximSection cfg,
  80. bool hasExchange,
  81. bool hasLoan,
  82. bool hasIntl,
  83. DateOnly day,
  84. CancellationToken ct
  85. ) {
  86. var touched = 0;
  87. if (hasExchange)
  88. {
  89. touched += await FetchExchangeAsync(db, client, baseUrl, cfg.ExchangeKey, day, ct);
  90. }
  91. if (hasLoan)
  92. {
  93. touched += await FetchRateAsync(db, client, baseUrl, cfg.LoanRateKey, "AP02", RateType.Loan, day, ct);
  94. }
  95. if (hasIntl)
  96. {
  97. touched += await FetchRateAsync(db, client, baseUrl, cfg.IntlRateKey, "AP03", RateType.International, day, ct);
  98. }
  99. if (touched > 0)
  100. {
  101. await db.SaveChangesAsync(ct);
  102. }
  103. else
  104. {
  105. Logger.LogInformation("[{Job}] searchdate={Day} 거시 미반영 (0건)", JobName, day);
  106. }
  107. }
  108. /// <summary>AP01 환율 수집 + upsert. 반환값 = 이번 호출로 add/update 한 행 수.</summary>
  109. private async Task<int> FetchExchangeAsync(IAppDbContext db, HttpClient client, string baseUrl, string authKey, DateOnly day, CancellationToken ct)
  110. {
  111. var url = $"{baseUrl}?authkey={authKey}&searchdate={day:yyyyMMdd}&data=AP01";
  112. var json = await KoreaEximHttp.GetStringWithRetryAsync(client, url, Logger, ct);
  113. var parsed = KoreaEximExchangeParser.Parse(json);
  114. Logger.LogInformation("[{Job}] AP01 환율 searchdate={Day} rows={Rows}", JobName, day, parsed.Count);
  115. if (parsed.Count == 0)
  116. {
  117. return 0;
  118. }
  119. var existing = await db.ExchangeRate.Where(c => c.TradeDate == day).ToListAsync(ct);
  120. var existingByUnit = existing.ToDictionary(c => c.CurrencyUnit);
  121. var count = 0;
  122. foreach (var row in parsed)
  123. {
  124. if (existingByUnit.TryGetValue(row.CurrencyUnit, out var rate))
  125. {
  126. rate.Update(row.CurrencyName, row.DealBasRate, row.Ttb, row.Tts);
  127. }
  128. else
  129. {
  130. var created = ExchangeRate.Create(row.CurrencyUnit, row.CurrencyName, day, row.DealBasRate, row.Ttb, row.Tts);
  131. await db.ExchangeRate.AddAsync(created, ct);
  132. existingByUnit[row.CurrencyUnit] = created;
  133. }
  134. count++;
  135. }
  136. return count;
  137. }
  138. /// <summary>AP02/AP03 금리 수집 + upsert. 반환값 = 이번 호출로 add/update 한 행 수.</summary>
  139. private async Task<int> FetchRateAsync(IAppDbContext db, HttpClient client, string baseUrl, string authKey, string data, RateType rateType, DateOnly day, CancellationToken ct)
  140. {
  141. var url = $"{baseUrl}?authkey={authKey}&searchdate={day:yyyyMMdd}&data={data}";
  142. var json = await KoreaEximHttp.GetStringWithRetryAsync(client, url, Logger, ct);
  143. var parsed = KoreaEximRateParser.Parse(json);
  144. Logger.LogInformation("[{Job}] {Data} {Type} searchdate={Day} rows={Rows}", JobName, data, rateType, day, parsed.Count);
  145. if (parsed.Count == 0)
  146. {
  147. return 0;
  148. }
  149. var existing = await db.InterestRate.Where(c => c.RateType == rateType && c.TradeDate == day).ToListAsync(ct);
  150. var existingByName = existing.ToDictionary(c => c.ItemName);
  151. var count = 0;
  152. foreach (var row in parsed)
  153. {
  154. if (existingByName.TryGetValue(row.ItemName, out var rate))
  155. {
  156. rate.Update(row.Rate);
  157. }
  158. else
  159. {
  160. var created = InterestRate.Create(rateType, row.ItemName, day, row.Rate);
  161. await db.InterestRate.AddAsync(created, ct);
  162. existingByName[row.ItemName] = created;
  163. }
  164. count++;
  165. }
  166. return count;
  167. }
  168. }