1772045028
2026-02-24 08:00:00
2026 年 2 月 24 日
スタニスワフ・ショコウスキー著
タグ:
私の中で 前の投稿Optimizely Commerce 用のメモリ効率の高いカタログ トラバーサル サービスを構築する方法を説明しました。このサービスはストリーミングを使用して、すべてを一度にメモリにロードせずに大規模なカタログを処理します。
しかし、適切に設計されたサービスを用意することは、戦いの半分に過ぎません。本当の価値は、運用シナリオで効果的に使用する方法を知ることから生まれます。この投稿では、カタログ データを処理するスケジュールされたジョブの実践的なパターンを、エラー処理、進捗レポート、回復戦略を含めて説明します。
簡単な要約: サービス
念のために言っておきますが、私たちが使用しているインターフェイスは次のとおりです。
public interface ICatalogTraversalService
{
IEnumerableICatalogTraversalItem> GetAllProducts(
CatalogTraversalOptions options,
CancellationToken cancellationToken = default);
}
public class CatalogTraversalOptions
{
public string? CatalogName { get; set; }
public ContentReference? CatalogLink { get; set; }
public DateTime? LastUpdated { get; set; }
}
このサービスは、カタログ階層を横断する際に、製品とバリエーションを一度に 1 つずつ生成します。では、効果的な使い方を見ていきましょう。
パターン 1: フルカタログのエクスポート
最も単純な使用例: カタログからすべての製品を外部システムにエクスポートします。このパターンは、初期データのロードまたは完全なリフレッシュに役立ちます。
[ScheduledPlugIn(
DisplayName = "[Catalog Traversal Demo] Export Catalog Products",
Description = "Exports all products from the Fashion catalog to external system",
GUID = "681EA6C4-B635-4CC3-8D9B-DBE3BEC602A6")]
public class CatalogExportJob : ScheduledJobBase
{
private readonly ICatalogTraversalService _catalogTraversal;
private readonly IExternalSystemClient _externalClient;
private readonly ILogger _logger;
private bool _stopSignaled;
public CatalogExportJob(
ICatalogTraversalService catalogTraversal,
IExternalSystemClient externalClient,
ILogger logger)
{
_catalogTraversal = catalogTraversal;
_externalClient = externalClient;
_logger = logger;
IsStoppable = true;
}
public override void Stop() => _stopSignaled = true;
public override string Execute()
{
var processedCount = 0;
var errorCount = 0;
var startTime = DateTime.UtcNow;
try
{
var options = new CatalogTraversalOptions
{
CatalogName = "Fashion"
};
_logger.LogInformation("Starting catalog export for '{CatalogName}'", options.CatalogName);
// The magic happens here - items are streamed one at a time
foreach (var item in _catalogTraversal.GetAllProducts(options, CancellationToken.None))
{
try
{
// Process each item - only one in memory at a time
switch (item)
{
case ProductContent product:
_externalClient.ExportProduct(product);
break;
case VariationContent variant:
_externalClient.ExportVariant(variant);
break;
}
processedCount++;
// Report progress every 100 items
if (processedCount % 100 == 0)
{
OnStatusChanged($"Processed {processedCount} items...");
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Error exporting item");
errorCount++;
// Continue processing despite errors
// Alternatively, you could fail fast by re-throwing
}
// Check if job was stopped by user
if (_stopSignaled)
{
_logger.LogWarning("Job stopped by user at {ProcessedCount} items", processedCount);
return $"Job stopped by user. Processed {processedCount} items.";
}
}
var duration = DateTime.UtcNow - startTime;
var result = $"Successfully processed {processedCount} items in {duration.TotalMinutes:F1} minutes. Errors: {errorCount}";
_logger.LogInformation("Catalog export completed: {Result}", result);
return result;
}
catch (Exception ex)
{
_logger.LogError(ex, "Fatal error in catalog export job");
return $"Job failed after processing {processedCount} items: {ex.Message}";
}
}
}
重要なポイント:
- 100 項目ごとの進捗レポートにより管理者に最新の情報が提供されます
- 個々のアイテムのエラーはログに記録されますが、ジョブ全体は停止しません
- 手動キャンセルの停止信号を尊重します
- 成功数とエラー数の両方を追跡して可視化します
パターン 2: 状態管理による増分同期
継続的な同期の場合は、最後に正常に実行されてから変更された項目のみを処理する必要があります。このパターンにより、処理時間と外部 API 呼び出しが大幅に削減されます。
[ScheduledPlugIn(
DisplayName = "[Catalog Traversal Demo] Incremental Catalog Sync",
Description = "Syncs only changed products since last run",
GUID = "A3BED4B6-FF3F-409E-895F-05567C8D3225")]
public class IncrementalCatalogSyncJob : ScheduledJobBase
{
private readonly ICatalogTraversalService _catalogTraversal;
private readonly ILastSyncRepository _lastSyncRepository;
private readonly IExternalSystemClient _externalClient;
private readonly ILogger _logger;
private const string SyncStateKey = "CatalogSync_Fashion";
public IncrementalCatalogSyncJob(
ICatalogTraversalService catalogTraversal,
ILastSyncRepository lastSyncRepository,
IExternalSystemClient externalClient,
ILogger logger)
{
_catalogTraversal = catalogTraversal;
_lastSyncRepository = lastSyncRepository;
_externalClient = externalClient;
_logger = logger;
}
public override string Execute()
{
// Get the last successful sync timestamp
var lastSyncDate = _lastSyncRepository.GetLastSyncDate(SyncStateKey);
var currentSyncDate = DateTime.UtcNow;
var updatedCount = 0;
var errorCount = 0;
try
{
var options = new CatalogTraversalOptions
{
CatalogName = "Fashion",
// Only get items updated since last sync
LastUpdated = lastSyncDate
};
var syncType = lastSyncDate.HasValue ? "Incremental" : "Full";
_logger.LogInformation(
"{SyncType} sync started. Last sync: {LastSyncDate}",
syncType,
lastSyncDate?.ToString("g") ?? "Never");
foreach (var item in _catalogTraversal.GetAllProducts(options, CancellationToken.None))
{
try
{
switch (item)
{
case ProductContent product:
_externalClient.SyncProduct(product);
break;
case VariationContent variant:
_externalClient.SyncVariant(variant);
break;
}
updatedCount++;
if (updatedCount % 50 == 0)
{
OnStatusChanged($"Synced {updatedCount} changed items...");
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Error syncing item");
errorCount++;
}
}
// Only update the last sync date if job completed successfully
_lastSyncRepository.SaveLastSyncDate(SyncStateKey, currentSyncDate);
var result = lastSyncDate.HasValue
? $"Incremental sync complete: {updatedCount} items changed since {lastSyncDate:g}. Errors: {errorCount}"
: $"Full sync complete: {updatedCount} items processed. Errors: {errorCount}";
_logger.LogInformation("Sync completed: {Result}", result);
return result;
}
catch (Exception ex)
{
// Don't update last sync date on failure - we'll retry from the same point next time
_logger.LogError(ex, "Fatal error in catalog sync job");
return $"Job failed after processing {updatedCount} items: {ex.Message}";
}
}
}
重要なポイント:
- 状態は正常に完了した場合にのみ保存されます
- 最初の実行では完全同期が実行されます (最終同期日はありません)
- 後続の実行は増分的ではるかに高速です
- 失敗した実行では状態が更新されないため、データが欠落していないことが保証されます
単純な状態リポジトリの実装
ここでは、同期状態を保存するために DDS を使用した基本的な実装を示します。
public interface ILastSyncRepository
{
DateTime? GetLastSyncDate(string key);
void SaveLastSyncDate(string key, DateTime date);
}
[EPiServerDataStore(AutomaticallyCreateStore = true, AutomaticallyRemapStore = true)]
public class SyncStateRecord : IDynamicData
{
public Identity Id { get; set; }
public string Key { get; set; }
public DateTime LastSyncDate { get; set; }
}
public class LastSyncRepository : ILastSyncRepository
{
private readonly DynamicDataStoreFactory _dataStoreFactory;
public LastSyncRepository(DynamicDataStoreFactory dataStoreFactory)
{
_dataStoreFactory = dataStoreFactory;
}
public DateTime? GetLastSyncDate(string key)
{
var store = _dataStoreFactory.GetStore(typeof(SyncStateRecord));
var record = store.FindSyncStateRecord>("Key", key).FirstOrDefault();
return record?.LastSyncDate;
}
public void SaveLastSyncDate(string key, DateTime date)
{
var store = _dataStoreFactory.GetStore(typeof(SyncStateRecord));
var record = store.FindSyncStateRecord>("Key", key).FirstOrDefault();
if (record == null)
{
record = new SyncStateRecord { Key = key, LastSyncDate = date };
store.Save(record);
}
else
{
record.LastSyncDate = date;
store.Save(record);
}
}
}
パターン 3: エラー回復を伴うバッチ処理
外部 API を呼び出すときは、多くの場合、効率を高めるためにリクエストをバッチ処理する必要があります。このパターンは、エラー耐性を維持しながらアイテムをバッチ処理する方法を示しています。
[ScheduledPlugIn(
DisplayName = "[Catalog Traversal Demo] Batch Catalog Export",
Description = "Exports products in batches to external API",
GUID = "C004BA30-C72F-445A-9EF6-EA3FCCF191B7")]
public class BatchCatalogExportJob : ScheduledJobBase
{
private readonly ICatalogTraversalService _catalogTraversal;
private readonly IExternalBatchClient _batchClient;
private readonly ILogger _logger;
private bool _stopSignaled;
private const int BatchSize = 50;
public BatchCatalogExportJob(
ICatalogTraversalService catalogTraversal,
IExternalBatchClient batchClient,
ILogger logger)
{
_catalogTraversal = catalogTraversal;
_batchClient = batchClient;
_logger = logger;
IsStoppable = true;
}
public override void Stop() => _stopSignaled = true;
public override string Execute()
{
var totalProcessed = 0;
var batchCount = 0;
var errorCount = 0;
try
{
var options = new CatalogTraversalOptions
{
CatalogName = "Fashion"
};
var batch = new List(BatchSize);
foreach (var item in _catalogTraversal.GetAllProducts(options, CancellationToken.None))
{
batch.Add(item);
// When batch is full, send it
if (batch.Count >= BatchSize)
{
var result = ProcessBatch(batch, ++batchCount);
totalProcessed += result.Processed;
errorCount += result.Errors;
batch.Clear();
OnStatusChanged($"Processed {totalProcessed} items in {batchCount} batches...");
}
if (_stopSignaled)
{
// Process remaining items before stopping
if (batch.Count > 0)
{
var result = ProcessBatch(batch, ++batchCount);
totalProcessed += result.Processed;
errorCount += result.Errors;
}
return $"Job stopped. Processed {totalProcessed} items in {batchCount} batches.";
}
}
// Process any remaining items in the last batch
if (batch.Count > 0)
{
var result = ProcessBatch(batch, ++batchCount);
totalProcessed += result.Processed;
errorCount += result.Errors;
}
return $"Successfully processed {totalProcessed} items in {batchCount} batches. Errors: {errorCount}";
}
catch (Exception ex)
{
_logger.LogError(ex, "Fatal error in batch export job");
return $"Job failed after processing {totalProcessed} items: {ex.Message}";
}
}
private (int Processed, int Errors) ProcessBatch(List batch, int batchNumber)
{
try
{
_logger.LogInformation("Processing batch {BatchNumber} with {ItemCount} items", batchNumber, batch.Count);
_batchClient.ExportBatch(batch);
return (batch.Count, 0);
}
catch (Exception ex)
{
_logger.LogError(ex, "Error processing batch {BatchNumber}", batchNumber);
// Fallback: try to process items individually
return ProcessBatchIndividually(batch, batchNumber);
}
}
private (int Processed, int Errors) ProcessBatchIndividually(List batch, int batchNumber)
{
_logger.LogWarning("Batch {BatchNumber} failed, attempting individual processing", batchNumber);
var processed = 0;
var errors = 0;
foreach (var item in batch)
{
try
{
_batchClient.ExportSingle(item);
processed++;
}
catch (Exception ex)
{
_logger.LogError(ex, "Error processing individual item from batch {BatchNumber}", batchNumber);
errors++;
}
}
return (processed, errors);
}
}
重要なポイント:
- アイテムは 50 個 (またはお好みのサイズ) のバッチにまとめられます。
- バッチが失敗した場合は、個別の処理にフォールバックします
- 停止中でも残りの項目は処理されます
- バッチ処理ロジックと個別処理ロジックを明確に分離
パターン 4: 複数カタログ処理
複数のカタログがあり、それらすべてを処理する必要がある場合、このパターンを使用すると、明確な分離と進行状況の可視化が保証されます。
[ScheduledPlugIn(
DisplayName = "[Catalog Traversal Demo] Multi-Catalog Sync",
Description = "Syncs all catalogs to external system",
GUID = "3FD79D80-47D3-4E17-85C1-5D97E196D691")]
public class MultiCatalogSyncJob : ScheduledJobBase
{
private readonly ICatalogTraversalService _catalogTraversal;
private readonly IContentLoader _contentLoader;
private readonly ReferenceConverter _referenceConverter;
private readonly IExternalSystemClient _externalClient;
private readonly ILogger _logger;
private bool _stopSignaled;
public MultiCatalogSyncJob(
ICatalogTraversalService catalogTraversal,
IContentLoader contentLoader,
ReferenceConverter referenceConverter,
IExternalSystemClient externalClient,
ILogger logger)
{
_catalogTraversal = catalogTraversal;
_contentLoader = contentLoader;
_referenceConverter = referenceConverter;
_externalClient = externalClient;
_logger = logger;
IsStoppable = true;
}
public override void Stop() => _stopSignaled = true;
public override string Execute()
{
var catalogResults = new Dictionary();
var totalProcessed = 0;
var totalErrors = 0;
try
{
// Get all catalogs
var catalogs = _contentLoader
.GetChildren(_referenceConverter.GetRootLink())
.ToList();
_logger.LogInformation("Found {CatalogCount} catalogs to process", catalogs.Count);
foreach (var catalog in catalogs)
{
if (_stopSignaled)
{
_logger.LogWarning("Job stopped while processing catalog '{CatalogName}'", catalog.Name);
break;
}
OnStatusChanged($"Processing catalog: {catalog.Name}");
var result = ProcessCatalog(catalog);
catalogResults[catalog.Name] = result;
totalProcessed += result.Processed;
totalErrors += result.Errors;
_logger.LogInformation(
"Completed catalog '{CatalogName}': {Processed} processed, {Errors} errors",
catalog.Name,
result.Processed,
result.Errors);
}
// Build detailed summary
var summary = new StringBuilder();
summary.AppendLine($"Multi-catalog sync completed:");
summary.AppendLine($"Total: {totalProcessed} items processed, {totalErrors} errors");
summary.AppendLine();
foreach (var (catalogName, result) in catalogResults)
{
summary.AppendLine($" {catalogName}: {result.Processed} items, {result.Errors} errors");
}
return summary.ToString();
}
catch (Exception ex)
{
_logger.LogError(ex, "Fatal error in multi-catalog sync job");
return $"Job failed. Processed {totalProcessed} items across {catalogResults.Count} catalogs.";
}
}
private (int Processed, int Errors) ProcessCatalog(CatalogContentBase catalog)
{
var processed = 0;
var errors = 0;
var options = new CatalogTraversalOptions
{
CatalogLink = catalog.ContentLink
};
foreach (var item in _catalogTraversal.GetAllProducts(options, CancellationToken.None))
{
try
{
switch (item)
{
case ProductContent product:
_externalClient.SyncProduct(product);
break;
case VariationContent variant:
_externalClient.SyncVariant(variant);
break;
}
processed++;
}
catch (Exception ex)
{
_logger.LogError(ex, "Error processing item in catalog '{CatalogName}'", catalog.Name);
errors++;
}
}
return (processed, errors);
}
}
重要なポイント:
- 各カタログは個別に処理されます
- 結果はカタログごとに追跡され、詳細なレポートが作成されます。
- ジョブはカタログ間で停止可能
- 概要には各カタログの結果が個別に表示されます
パターン 5: 進捗状況の追跡と監視
長時間実行されるジョブの場合、詳細な進行状況の追跡により、運用チームがパフォーマンスを監視し、問題を早期に特定するのに役立ちます。
[ScheduledPlugIn(
DisplayName = "[Catalog Traversal Demo] Catalog Export with Detailed Progress",
Description = "Exports catalog with detailed progress tracking and metrics",
GUID = "20DB2FB3-0944-4B19-BE73-AD2E17E9FED0")]
public class DetailedProgressExportJob : ScheduledJobBase
{
private readonly ICatalogTraversalService _catalogTraversal;
private readonly IExternalSystemClient _externalClient;
private readonly ILogger _logger;
private bool _stopSignaled;
public DetailedProgressExportJob(
ICatalogTraversalService catalogTraversal,
IExternalSystemClient externalClient,
ILogger logger)
{
_catalogTraversal = catalogTraversal;
_externalClient = externalClient;
_logger = logger;
IsStoppable = true;
}
public override void Stop() => _stopSignaled = true;
public override string Execute()
{
var metrics = new ProcessingMetrics();
var progressReporter = new ProgressReporter(this, _logger);
try
{
var options = new CatalogTraversalOptions
{
CatalogName = "Fashion"
};
foreach (var item in _catalogTraversal.GetAllProducts(options, CancellationToken.None))
{
try
{
var processingStarted = DateTime.UtcNow;
switch (item)
{
case ProductContent product:
_externalClient.ExportProduct(product);
metrics.ProductsProcessed++;
break;
case VariationContent variant:
_externalClient.ExportVariant(variant);
metrics.VariantsProcessed++;
break;
}
metrics.RecordProcessingTime(DateTime.UtcNow - processingStarted);
}
catch (Exception ex)
{
_logger.LogError(ex, "Error processing item");
metrics.Errors++;
}
// Report progress with detailed metrics
progressReporter.ReportProgress(metrics);
if (_stopSignaled)
{
return metrics.GetStoppedSummary();
}
}
return metrics.GetCompletedSummary();
}
catch (Exception ex)
{
_logger.LogError(ex, "Fatal error in export job");
return metrics.GetFailedSummary(ex.Message);
}
}
private class ProcessingMetrics
{
public int ProductsProcessed { get; set; }
public int VariantsProcessed { get; set; }
public int Errors { get; set; }
public int TotalProcessed => ProductsProcessed + VariantsProcessed;
public DateTime StartTime { get; } = DateTime.UtcNow;
private readonly List _processingTimes = new();
public void RecordProcessingTime(TimeSpan time)
{
_processingTimes.Add(time);
// Keep only last 100 samples to calculate average
if (_processingTimes.Count > 100)
{
_processingTimes.RemoveAt(0);
}
}
public TimeSpan AverageProcessingTime =>
_processingTimes.Any()
? TimeSpan.FromTicks((long)_processingTimes.Average(t => t.Ticks))
: TimeSpan.Zero;
public TimeSpan ElapsedTime => DateTime.UtcNow - StartTime;
public double ItemsPerSecond =>
ElapsedTime.TotalSeconds > 0
? TotalProcessed / ElapsedTime.TotalSeconds
: 0;
public string GetCompletedSummary()
{
return $@"Export completed successfully:
Products: {ProductsProcessed}
Variants: {VariantsProcessed}
Total: {TotalProcessed}
Errors: {Errors}
Duration: {ElapsedTime.TotalMinutes:F1} minutes
Average: {ItemsPerSecond:F1} items/second
Avg processing time: {AverageProcessingTime.TotalMilliseconds:F0}ms";
}
public string GetStoppedSummary()
{
return $"Job stopped. Processed {TotalProcessed} items ({ProductsProcessed} products, {VariantsProcessed} variants). Errors: {Errors}";
}
public string GetFailedSummary(string error)
{
return $"Job failed after processing {TotalProcessed} items: {error}";
}
}
private class ProgressReporter
{
private readonly DetailedProgressExportJob _job;
private readonly ILogger _logger;
private DateTime _lastReport = DateTime.MinValue;
private int _lastReportedCount = 0;
private const int ReportIntervalSeconds = 10;
public ProgressReporter(DetailedProgressExportJob job, ILogger logger)
{
_job = job;
_logger = logger;
}
public void ReportProgress(ProcessingMetrics metrics)
{
var now = DateTime.UtcNow;
// Report every 10 seconds
if ((now - _lastReport).TotalSeconds 0
? itemsSinceLastReport / timeSinceLastReport.TotalSeconds
: 0;
var status = $@"Progress: {metrics.TotalProcessed} items ({metrics.ProductsProcessed}p/{metrics.VariantsProcessed}v)
Rate: {metrics.ItemsPerSecond:F1} items/s (recent: {recentRate:F1} items/s)
Errors: {metrics.Errors}
Elapsed: {metrics.ElapsedTime.TotalMinutes:F1}m";
_job.OnStatusChanged(status);
_logger.LogInformation("Job progress: {Status}", status.Replace("n", " | "));
_lastReport = now;
_lastReportedCount = metrics.TotalProcessed;
}
}
}
重要なポイント:
- 詳細な指標を追跡: 製品とバリアント、処理時間、スループット
- 現在および最近のレートで 10 秒ごとに進捗状況をレポートします
- パフォーマンス監視の平均処理時間を計算します。
- 完了、停止、または失敗に関する詳細な概要を提供します
各パターンをいつ使用するか
ニーズに基づいて適切なパターンを選択してください。
| パターン | 最適な用途 | 主な利点 |
|---|---|---|
| 完全なエクスポート | 初期ロード、完全なリフレッシュ | シンプルでわかりやすい |
| 増分同期 | 継続的な同期 | 以降の実行が大幅に高速化 |
| バッチ処理 | API レート制限、効率 | API 呼び出しを減らし、障害を適切に処理します |
| マルチカタログ | 複数のカタログ、複雑なセットアップ | きれいな分離、詳細なレポート作成 |
| 詳細な進捗状況 | 長時間実行されるジョブ、監視 | 運用の可視性、パフォーマンスの洞察 |
パターンを組み合わせることもできます。たとえば、最も効率的に継続的な同期を行うには、増分同期とバッチ処理を使用します。
ベストプラクティス
上記のパターンに基づいて、いくつかの重要な推奨事項を次に示します。
エラー処理:
- 個々のアイテムのエラーをログに記録しますが、処理は続行します
- フォールバック戦略の実装 (バッチ → 個別)
- 本当に致命的なエラーの場合のみフェイルファストします
状態管理:
- 正常に完了した後にのみ同期状態を保存する
- さまざまな同期ジョブに一意のキーを使用する
- 追加のメタデータ (項目数、期間) を保存することを検討してください。
進捗報告:
- 進捗状況を定期的に報告します (10 秒ごと、または 100 項目ごと)
- 意味のある指標 (アイテム/秒、エラー、経過時間) を含める
- 全体的なパフォーマンスと最近のパフォーマンスの両方を表示する
キャンセル:
- 常に停止信号を尊重してください
- 停止する前に残りのアイテム/バッチを処理します
- 完了した内容について明確なステータスを提供する
ロギング:
- 意味のあるコンテキストを含む構造化ログを使用する
- 適切なレベルでログを記録します (進行状況に関する情報、問題に関する警告)
- トラブルシューティング用の識別子(カタログ名、品目コード)を含める
まとめ
カタログ トラバーサル サービスは強固な基盤を提供しますが、本当の価値は、スケジュールされたジョブで効果的に使用することで得られます。この投稿のパターンは、最も一般的なシナリオをカバーしています。
- 完全なデータロードのための完全なエクスポート
- 増分同期による効率的な継続的な更新
- API の効率性と復元力を高めるバッチ処理
- 複雑なセットアップのためのマルチカタログ処理
- 運用状況を可視化するための詳細な進捗状況の追跡
ニーズに合ったパターンを選択し、より洗練されたシナリオのために躊躇せずに組み合わせてください。ストリーミング アプローチにより、カタログの成長に合わせてジョブが適切に拡張されます。
他にも取り上げてほしいパターンやユースケースはありますか?コメントで知らせてください!
読んでいただきありがとうございます。これらのパターンが、Optimizely Commerce ソリューションで堅牢なカタログ処理ジョブを構築するのに役立つことを願っています。
この投稿はシリーズの一部です
- パート 1: サービスの構築
- パート 2: 実際のスケジュールされたジョブ パターン – (この投稿)
- パート 3: Hangfire の統合 – 将来のリリースをお待ちください!
#カタログ走査の実行中パート #実際のスケジュールされたジョブ #パターン