Обход каталога в действии. Часть 2. Реальные шаблоны запланированных заданий

24 февраля 2026 г.

Станислав Шолковский

Теги:


.СЕТЬ
(2)


каталог
(2)


коммерция
(2)


эписервер
(10)


оптимизированный
(10)


узоры
(1)


запланированные задания
(3)

В моем предыдущий постЯ показал, как создать службу обхода каталога с эффективным использованием памяти для 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}";
        }
    }
}

  

×


[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}";
        }
    }
}

  

Read more:  Брат Мики Парсонс раскрывает «отвратительную часть» о торговле ковбоями

×


[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);
    }
}
  

×


[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 штук (или по вашему выбору).
  • Если пакет не выполнен, вернитесь к индивидуальной обработке.
  • Остальные элементы обрабатываются даже после остановки
  • Четкое разделение между пакетной и индивидуальной логикой обработки
Read more:  «Николандрия» - это не просто культурное явление, это любимая часть истории реальности пары

Схема 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);
    }
}
  

×


[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;
        }
    }
}

  

Read more:  Республиканцы Калифорнии просят Верховный суд США заблокировать новую карту Конгресса | Калифорния

×


[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 — ждите будущего выпуска!

Читайте также

Leave a Comment

This site uses Akismet to reduce spam. Learn how your comment data is processed.