WorldIndexSyncService.cs 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124
  1. using System.Text.Json;
  2. using Application.Abstractions.Data;
  3. using Domain.Entities.Stocks;
  4. using Microsoft.EntityFrameworkCore;
  5. using Microsoft.Extensions.DependencyInjection;
  6. using Microsoft.Extensions.Logging;
  7. using Microsoft.Extensions.Options;
  8. using SharedKernel;
  9. namespace Infrastructure.StockData;
  10. /// <summary>
  11. /// 세계 주요국 증시지수 스냅샷 수집 — Yahoo Finance v8 chart(무키)를 심볼별로 하루 1회 조회해 최신 종가를 upsert.
  12. /// config(WorldIndex:Items) 심볼(Yahoo 표기, 예: ^GSPC/^N225/000001.SS)을 하나씩 GET 하고 meta 를 파싱해 WorldIndexSnapshot 에 반영한다.
  13. /// 국내(코스피)는 이 배치가 아니라 GetWorldIndices Handler 가 IndexDailyPrice 실데이터에서 병합한다.
  14. /// 기본 07:00 KST 실행. 심볼 하나 실패는 skip(격리)하고, 전량 실패 시에만 재시도. (Stooq 봇차단으로 소스 전환 2026-07-09)
  15. /// </summary>
  16. internal sealed class WorldIndexSyncService(
  17. IServiceScopeFactory scopeFactory,
  18. IHttpClientFactory httpClientFactory,
  19. IOptions<AppSettings> settings,
  20. ILogger<WorldIndexSyncService> logger
  21. ) : DailyScheduledService(logger)
  22. {
  23. protected override string JobName => "WorldIndexSync";
  24. protected override TimeOnly TargetTime => ParseTime(settings.Value.WorldIndex.SyncTime, new TimeOnly(7, 0));
  25. protected override int MaxRetryCount => 2;
  26. protected override TimeSpan RetryDelay => TimeSpan.FromMinutes(settings.Value.WorldIndex.RetryDelayMinutes > 0 ? settings.Value.WorldIndex.RetryDelayMinutes : 60);
  27. protected override async Task<bool> RunOnceAsync(DateOnly todayKst, CancellationToken ct)
  28. {
  29. var cfg = settings.Value.WorldIndex;
  30. using var scope = scopeFactory.CreateScope();
  31. var collectorSettings = scope.ServiceProvider.GetRequiredService<ICollectorSettingsProvider>();
  32. if (!await collectorSettings.IsEnabledAsync(CollectorFlag.WorldIndex, ct))
  33. {
  34. return true;
  35. }
  36. var items = cfg.Items.Where(c => !string.IsNullOrWhiteSpace(c.Symbol)).ToList();
  37. if (items.Count == 0)
  38. {
  39. Logger.LogWarning("[{Job}] WorldIndex:Items 미설정 — 수집 skip", JobName);
  40. return true;
  41. }
  42. var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
  43. var client = httpClientFactory.CreateClient(YahooFinanceHttp.ClientName);
  44. var existing = await db.WorldIndexSnapshot.ToListAsync(ct);
  45. var existingBySymbol = existing.ToDictionary(c => c.Symbol, StringComparer.OrdinalIgnoreCase);
  46. var baseUrl = cfg.BaseUrl.TrimEnd('/');
  47. var inserted = 0;
  48. var updated = 0;
  49. var failed = 0;
  50. foreach (var meta in items)
  51. {
  52. ct.ThrowIfCancellationRequested();
  53. var symbol = meta.Symbol.Trim();
  54. // range=1mo — 최신 종가(meta)와 함께 스파크라인용 일별 종가 시계열(~20영업일)을 한 번에 수신
  55. var url = $"{baseUrl}/v8/finance/chart/{Uri.EscapeDataString(symbol)}?interval=1d&range=1mo";
  56. YahooChartParser.YahooQuote? q;
  57. string? sparkJson = null;
  58. try
  59. {
  60. var json = await YahooFinanceHttp.GetStringWithRetryAsync(client, url, Logger, ct);
  61. q = YahooChartParser.Parse(json, symbol);
  62. var series = YahooChartParser.ParseCloseSeries(json, 30);
  63. if (series.Count > 0)
  64. {
  65. sparkJson = JsonSerializer.Serialize(series);
  66. }
  67. }
  68. catch (Exception ex) when (ex is not OperationCanceledException)
  69. {
  70. Logger.LogWarning(ex, "[{Job}] {Symbol} 조회 실패 — skip", JobName, symbol);
  71. failed++;
  72. continue;
  73. }
  74. if (q is null)
  75. {
  76. Logger.LogWarning("[{Job}] {Symbol} 데이터 없음 — skip", JobName, symbol);
  77. failed++;
  78. continue;
  79. }
  80. if (existingBySymbol.TryGetValue(q.Symbol, out var snap))
  81. {
  82. snap.UpdateMeta(meta.Name, meta.CountryCode, meta.ExchangeName);
  83. snap.Apply(q.TradeDate, q.Close, q.Open, q.High, q.Low);
  84. snap.SetSpark(sparkJson);
  85. updated++;
  86. }
  87. else
  88. {
  89. var created = WorldIndexSnapshot.Create(q.Symbol, meta.Name, meta.CountryCode, meta.ExchangeName, q.TradeDate, q.Close, q.Open, q.High, q.Low);
  90. created.SetSpark(sparkJson);
  91. await db.WorldIndexSnapshot.AddAsync(created, ct);
  92. existingBySymbol[q.Symbol] = created;
  93. inserted++;
  94. }
  95. }
  96. if (inserted == 0 && updated == 0)
  97. {
  98. Logger.LogWarning("[{Job}] 전량 실패 (failed={Failed}) — 재시도 대상", JobName, failed);
  99. return false;
  100. }
  101. await db.SaveChangesAsync(ct);
  102. Logger.LogInformation("[{Job}] 완료 — inserted={Inserted}, updated={Updated}, failed={Failed}", JobName, inserted, updated, failed);
  103. return true;
  104. }
  105. }