| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210 |
- using Application.Abstractions.Data;
- using Application.Helpers;
- using Domain.Entities.Stocks;
- using Domain.Entities.Stocks.ValueObject;
- using Microsoft.EntityFrameworkCore;
- using Microsoft.Extensions.DependencyInjection;
- using Microsoft.Extensions.Logging;
- using Microsoft.Extensions.Options;
- using SharedKernel;
- namespace Infrastructure.StockData;
- /// <summary>
- /// 거시·환율 수집 (한국수출입은행 OpenAPI) — AP01 현재환율 + AP02 대출금리 + AP03 국제금리. 기본 11:30 KST 실행
- /// (koreaexim 은 영업일 ~11시 이후 반영). 공용 KrxBackfill 로 [today-MacroBackfillDays, 직전 영업일] 을 최신일부터 과거로 훑으며
- /// 미적재일만 채운다 — quota 보호를 위해 1회 실행당 BackfillMaxPerRun(기본 60)일까지만 fetch. 각 날짜는 세 data 종류를
- /// (키가 설정된 것만) 모두 수집해 ExchangeRate / InterestRate upsert (UQ = (CurrencyUnit,TradeDate) / (RateType,ItemName,TradeDate)).
- /// 빈 배열([]) — 주말/공휴일·오전 미반영 — 은 "데이터 없음" 으로 취급(그 날은 미적재 유지, 다음 실행 때 재시도). 키 전부 미설정 시 skip.
- /// existsForDate 판정은 세 종류 중 하나라도 그 날짜에 적재됐는지로 본다 (부분 성공 후 재실행 시 남은 종류를 다시 채울 수 있게 upsert 로 안전).
- /// </summary>
- internal sealed class MacroDataSyncService(
- IServiceScopeFactory scopeFactory,
- IHttpClientFactory httpClientFactory,
- IOptions<AppSettings> settings,
- ILogger<MacroDataSyncService> logger
- ) : DailyScheduledService(logger)
- {
- protected override string JobName => "MacroDataSync";
- protected override TimeOnly TargetTime => ParseTime(settings.Value.KoreaExim.MacroSyncTime, new TimeOnly(11, 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.KoreaExim;
- using var scope = scopeFactory.CreateScope();
- var collectorSettings = scope.ServiceProvider.GetRequiredService<ICollectorSettingsProvider>();
- if (!await collectorSettings.IsEnabledAsync(CollectorFlag.KoreaEximMacro, ct))
- {
- return true;
- }
- cfg = cfg with { ExchangeKey = await collectorSettings.GetKeyAsync(CollectorKey.KoreaEximExchange, ct) ?? cfg.ExchangeKey, LoanRateKey = await collectorSettings.GetKeyAsync(CollectorKey.KoreaEximLoanRate, ct) ?? cfg.LoanRateKey, IntlRateKey = await collectorSettings.GetKeyAsync(CollectorKey.KoreaEximIntlRate, ct) ?? cfg.IntlRateKey };
- var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
- var hasExchange = !string.IsNullOrWhiteSpace(cfg.ExchangeKey);
- var hasLoan = !string.IsNullOrWhiteSpace(cfg.LoanRateKey);
- var hasIntl = !string.IsNullOrWhiteSpace(cfg.IntlRateKey);
- if (!hasExchange && !hasLoan && !hasIntl)
- {
- Logger.LogWarning("[{Job}] KoreaExim 키(ExchangeKey/LoanRateKey/IntlRateKey) 전부 미설정 — 수집 skip", JobName);
- return true;
- }
- var client = httpClientFactory.CreateClient(KoreaEximHttp.ClientName);
- var endDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, todayKst.AddDays(1), ct);
- var days = cfg.MacroBackfillDays > 0 ? cfg.MacroBackfillDays : 30;
- var startDate = todayKst.AddDays(-days);
- var holidays = (await db.MarketHoliday.AsNoTracking()
- .Where(c => c.Date >= startDate && c.Date <= endDate)
- .Select(c => c.Date)
- .ToListAsync(ct)).ToHashSet();
- var maxPerRun = cfg.BackfillMaxPerRun > 0 ? cfg.BackfillMaxPerRun : 60;
- var baseUrl = cfg.BaseUrl.TrimEnd('/');
- var fetched = await KrxBackfill.RunAsync(
- existsForDate: (day, token) => ExistsForDateAsync(db, day, hasExchange, token),
- fetchAndUpsertForDate: (day, token) => FetchAndUpsertAsync(db, client, baseUrl, cfg, hasExchange, hasLoan, hasIntl, day, token),
- startDate: startDate,
- endDate: endDate,
- holidays: holidays,
- maxPerRun: maxPerRun,
- delayMs: 300,
- ct: ct);
- Logger.LogInformation("[{Job}] 완료 — 창=[{Start}~{End}], 이번 실행 fetch={Fetched}일 (maxPerRun={Max})",
- JobName, startDate, endDate, fetched, maxPerRun);
- return true;
- }
- /// <summary>해당 날짜 적재 여부 — 환율 키가 있으면 ExchangeRate, 없으면 InterestRate 로 판정</summary>
- private static async Task<bool> ExistsForDateAsync(IAppDbContext db, DateOnly day, bool hasExchange, CancellationToken ct)
- {
- if (hasExchange)
- {
- return await db.ExchangeRate.AsNoTracking().AnyAsync(c => c.TradeDate == day, ct);
- }
- return await db.InterestRate.AsNoTracking().AnyAsync(c => c.TradeDate == day, ct);
- }
- private async Task FetchAndUpsertAsync(
- IAppDbContext db,
- HttpClient client,
- string baseUrl,
- AppSettings.KoreaEximSection cfg,
- bool hasExchange,
- bool hasLoan,
- bool hasIntl,
- DateOnly day,
- CancellationToken ct
- ) {
- var touched = 0;
- if (hasExchange)
- {
- touched += await FetchExchangeAsync(db, client, baseUrl, cfg.ExchangeKey, day, ct);
- }
- if (hasLoan)
- {
- touched += await FetchRateAsync(db, client, baseUrl, cfg.LoanRateKey, "AP02", RateType.Loan, day, ct);
- }
- if (hasIntl)
- {
- touched += await FetchRateAsync(db, client, baseUrl, cfg.IntlRateKey, "AP03", RateType.International, day, ct);
- }
- if (touched > 0)
- {
- await db.SaveChangesAsync(ct);
- }
- else
- {
- Logger.LogInformation("[{Job}] searchdate={Day} 거시 미반영 (0건)", JobName, day);
- }
- }
- /// <summary>AP01 환율 수집 + upsert. 반환값 = 이번 호출로 add/update 한 행 수.</summary>
- private async Task<int> FetchExchangeAsync(IAppDbContext db, HttpClient client, string baseUrl, string authKey, DateOnly day, CancellationToken ct)
- {
- var url = $"{baseUrl}?authkey={authKey}&searchdate={day:yyyyMMdd}&data=AP01";
- var json = await KoreaEximHttp.GetStringWithRetryAsync(client, url, Logger, ct);
- var parsed = KoreaEximExchangeParser.Parse(json);
- Logger.LogInformation("[{Job}] AP01 환율 searchdate={Day} rows={Rows}", JobName, day, parsed.Count);
- if (parsed.Count == 0)
- {
- return 0;
- }
- var existing = await db.ExchangeRate.Where(c => c.TradeDate == day).ToListAsync(ct);
- var existingByUnit = existing.ToDictionary(c => c.CurrencyUnit);
- var count = 0;
- foreach (var row in parsed)
- {
- if (existingByUnit.TryGetValue(row.CurrencyUnit, out var rate))
- {
- rate.Update(row.CurrencyName, row.DealBasRate, row.Ttb, row.Tts);
- }
- else
- {
- var created = ExchangeRate.Create(row.CurrencyUnit, row.CurrencyName, day, row.DealBasRate, row.Ttb, row.Tts);
- await db.ExchangeRate.AddAsync(created, ct);
- existingByUnit[row.CurrencyUnit] = created;
- }
- count++;
- }
- return count;
- }
- /// <summary>AP02/AP03 금리 수집 + upsert. 반환값 = 이번 호출로 add/update 한 행 수.</summary>
- private async Task<int> FetchRateAsync(IAppDbContext db, HttpClient client, string baseUrl, string authKey, string data, RateType rateType, DateOnly day, CancellationToken ct)
- {
- var url = $"{baseUrl}?authkey={authKey}&searchdate={day:yyyyMMdd}&data={data}";
- var json = await KoreaEximHttp.GetStringWithRetryAsync(client, url, Logger, ct);
- var parsed = KoreaEximRateParser.Parse(json);
- Logger.LogInformation("[{Job}] {Data} {Type} searchdate={Day} rows={Rows}", JobName, data, rateType, day, parsed.Count);
- if (parsed.Count == 0)
- {
- return 0;
- }
- var existing = await db.InterestRate.Where(c => c.RateType == rateType && c.TradeDate == day).ToListAsync(ct);
- var existingByName = existing.ToDictionary(c => c.ItemName);
- var count = 0;
- foreach (var row in parsed)
- {
- if (existingByName.TryGetValue(row.ItemName, out var rate))
- {
- rate.Update(row.Rate);
- }
- else
- {
- var created = InterestRate.Create(rateType, row.ItemName, day, row.Rate);
- await db.InterestRate.AddAsync(created, ct);
- existingByName[row.ItemName] = created;
- }
- count++;
- }
- return count;
- }
- }
|