24 февраля 2026 г.
Станислав Шолковский
Теги:
В моем предыдущий постЯ показал, как создать службу обхода каталога с эффективным использованием памяти для Optimizely Commerce. Сервис использует потоковую передачу для обработки больших каталогов без загрузки всего в память сразу.
Но иметь хорошо продуманный сервис – это только полдела. Реальная ценность заключается в знании того, как эффективно использовать его в производственных сценариях. В этом посте я расскажу о практических шаблонах запланированных заданий по обработке данных каталога, а также об обработке ошибок, составлении отчетов о ходе работы и стратегиях устойчивости.
Краткий обзор: Сервис
Напоминаем, вот интерфейс, с которым мы работаем:
public interface ICatalogTraversalService
{
IEnumerable<ICatalogTraversalItem> 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: полный экспорт каталога
Самый простой вариант использования: экспортировать все товары из каталога во внешнюю систему. Этот шаблон полезен для начальной загрузки данных или полного обновления.
[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.Find<SyncStateRecord>("Key", key).FirstOrDefault();
return record?.LastSyncDate;
}
public void SaveLastSyncDate(string key, DateTime date)
{
var store = _dataStoreFactory.GetStore(typeof(SyncStateRecord));
var record = store.Find<SyncStateRecord>("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 < ReportIntervalSeconds)
{
return;
}
var itemsSinceLastReport = metrics.TotalProcessed - _lastReportedCount;
var timeSinceLastReport = now - _lastReport;
var recentRate = timeSinceLastReport.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 — ждите будущего выпуска!
Читайте также

