diff --git a/.env.example b/.env.example index ae70d45..860ada6 100644 --- a/.env.example +++ b/.env.example @@ -44,8 +44,15 @@ Media__FfmpegPath=/usr/bin/ffmpeg Media__FfprobePath=/usr/bin/ffprobe # Максимальный размер загружаемого файла (20 ГБ). Media__MaxUploadBytes=21474836480 -# Потоки ffmpeg — оставляем ядро API/раздаче (сервер 4 vCPU). +# Потоки одного ffmpeg (0 — авто/все ядра). Media__TranscodeThreads=3 +# Сколько файлов транскодировать одновременно (1 — по умолчанию, как раньше). +# Тюнинг под ядра: обычно TranscodeThreads × MaxParallelTranscodes ≈ число ядер. +# 4 vCPU: Threads=3, Parallel=1 (или Threads=2, Parallel=2) +# 8 vCPU: Threads=4, Parallel=2 ← лучший throughput на пачках серий +# (или Threads=6, Parallel=1 — если важнее скорость одного файла) +# Память: один 1080p-транскод ~0.5 ГБ; Parallel=2 требует ~1 ГБ + база API. +Media__MaxParallelTranscodes=1 # Период опроса каталога inbox/ сканером, секунды. Media__InboxScanSeconds=15 # Нормализация громкости при обработке (EBU R128 loudnorm) — все ролики одинаковой громкости. diff --git a/backend/src/TeleWave.Infrastructure/Media/MediaOptions.cs b/backend/src/TeleWave.Infrastructure/Media/MediaOptions.cs index 46a922a..503dd98 100644 --- a/backend/src/TeleWave.Infrastructure/Media/MediaOptions.cs +++ b/backend/src/TeleWave.Infrastructure/Media/MediaOptions.cs @@ -10,9 +10,15 @@ public sealed class MediaOptions /// Максимальный размер загружаемого файла, байт (по умолчанию 20 ГБ). public long MaxUploadBytes { get; init; } = 20L * 1024 * 1024 * 1024; - /// Число потоков ffmpeg — оставляем ядро API и раздаче (см. docs, 4 vCPU). + /// Число потоков одного ffmpeg (0 — авто/все ядра). На многоядерных лучше сочетать с + /// : несколько файлов по меньшему числу потоков дают больший + /// суммарный throughput, чем один файл на всех ядрах. public int TranscodeThreads { get; init; } = 3; + /// Сколько файлов транскодировать одновременно (по умолчанию 1 — как было). Подбирается + /// под ядра: ориентир — ядра / TranscodeThreads, но не больше, чем позволяет память/диск. + public int MaxParallelTranscodes { get; init; } = 1; + /// Период опроса inbox/ сканером, секунды. public int InboxScanSeconds { get; init; } = 15; diff --git a/backend/src/TeleWave.Infrastructure/Media/MediaProcessingBackgroundService.cs b/backend/src/TeleWave.Infrastructure/Media/MediaProcessingBackgroundService.cs index 33d5e33..4f29969 100644 --- a/backend/src/TeleWave.Infrastructure/Media/MediaProcessingBackgroundService.cs +++ b/backend/src/TeleWave.Infrastructure/Media/MediaProcessingBackgroundService.cs @@ -1,3 +1,4 @@ +using System.Collections.Concurrent; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; @@ -9,11 +10,13 @@ using TeleWave.Domain.Media; namespace TeleWave.Infrastructure.Media; /// -/// Единственный обработчик медиа: по одному ассету за раз прогоняет через ffmpeg. Источник истины — -/// статус в БД: сервис в цикле берёт следующий Pending из базы и обрабатывает, а очередь лишь будит -/// его без задержки. Поэтому рестарт/краш ничего не теряет — незавершённые задачи подхватываются из -/// БД (прерванные Processing на старте сбрасываются в Pending). БД-контекст держится короткими -/// отрезками (пометить статус), сам транскод идёт вне scope, чтобы не держать соединение минутами. +/// Обработчик медиа: прогоняет ассеты через ffmpeg, до +/// файлов одновременно. Источник истины — статус в БД: единственный диспетчер последовательно и +/// атомарно захватывает следующий Pending (помечает Processing), поэтому два транскода никогда не +/// возьмут один ассет; сам транскод запускается в фоне с ограничением по числу слотов. Рестарт/краш +/// ничего не теряет — незавершённые подхватываются из БД (прерванные Processing на старте сбрасываются +/// в Pending). БД-контекст держится короткими отрезками (пометить статус), сам транскод идёт вне +/// scope, чтобы не держать соединение минутами. /// public sealed class MediaProcessingBackgroundService( IMediaProcessingQueue queue, @@ -21,25 +24,59 @@ public sealed class MediaProcessingBackgroundService( MediaPathResolver paths, IMediaProcessor processor, IOptions storageOptions, + IOptions mediaOptions, ILogger logger ) : BackgroundService { // Периодически перепроверяем БД, даже если сигнал не пришёл — страховка на любой случай. private static readonly TimeSpan IdlePoll = TimeSpan.FromSeconds(30); private readonly StorageOptions _storage = storageOptions.Value; + private readonly int _maxParallel = Math.Max(1, mediaOptions.Value.MaxParallelTranscodes); protected override async Task ExecuteAsync(CancellationToken stoppingToken) { paths.EnsureDirectories(); await ResetInterruptedAsync(stoppingToken); + // Слоты параллелизма: не запускаем больше _maxParallel транскодов одновременно. + using var slots = new SemaphoreSlim(_maxParallel, _maxParallel); + var inFlight = new ConcurrentDictionary(); + while (!stoppingToken.IsCancellationRequested) { try { - // Разобрать всю накопившуюся работу из БД. - while (await NextPendingIdAsync(stoppingToken) is { } assetId) - await ProcessAsync(assetId, stoppingToken); + // Захватываем и раздаём по слотам всю накопившуюся работу из БД. + while (!stoppingToken.IsCancellationRequested) + { + await slots.WaitAsync(stoppingToken); + var claim = await ClaimNextAsync(stoppingToken); + if (claim is not { } job) + { + slots.Release(); + break; + } + + var task = Task.Run( + async () => + { + try + { + await ProcessClaimedAsync(job.Id, job.Extension, stoppingToken); + } + finally + { + slots.Release(); + } + }, + CancellationToken.None + ); + inFlight[task] = 0; + _ = task.ContinueWith( + t => inFlight.TryRemove(t, out _), + TaskScheduler.Default + ); + } // Работы нет — ждём сигнала о новой либо периодического опроса. using var wake = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); @@ -63,6 +100,16 @@ public sealed class MediaProcessingBackgroundService( await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken); } } + + // Даём уже запущенным транскодам корректно завершиться (или отмениться) на остановке. + try + { + await Task.WhenAll(inFlight.Keys.ToArray()); + } + catch + { + // Ошибки/отмена отдельных задач уже залогированы внутри ProcessClaimedAsync. + } } /// Сброс прерванных рестартом задач (Processing → Pending) на старте. @@ -82,26 +129,31 @@ public sealed class MediaProcessingBackgroundService( await db.SaveChangesAsync(cancellationToken); } - /// Id самого раннего ассета в статусе Pending, либо null если работы нет. - private async Task NextPendingIdAsync(CancellationToken cancellationToken) + /// + /// Атомарно захватывает самый ранний Pending: помечает его Processing и возвращает (id, расширение), + /// либо null если работы нет. Вызывается только диспетчером последовательно, поэтому два транскода + /// не возьмут один ассет. + /// + private async Task<(Guid Id, string Extension)?> ClaimNextAsync(CancellationToken cancellationToken) { await using var scope = scopeFactory.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); - var id = await db.MediaAssets + var asset = await db.MediaAssets .Where(x => x.Status == MediaAssetStatus.Pending) .OrderBy(x => x.CreatedAt) - .Select(x => (Guid?)x.Id) .FirstOrDefaultAsync(cancellationToken); - return id; + if (asset is null) + return null; + + asset.MarkProcessing(); + await db.SaveChangesAsync(cancellationToken); + return (asset.Id, asset.OriginalExtension); } - private async Task ProcessAsync(Guid assetId, CancellationToken cancellationToken) + /// Обрабатывает уже захваченный (Processing) ассет: транскод → Ready/Failed. + private async Task ProcessClaimedAsync(Guid assetId, string extension, CancellationToken cancellationToken) { - var extension = await BeginProcessingAsync(assetId, cancellationToken); - if (extension is null) - return; - try { var result = await processor.ProcessAsync(assetId, extension, cancellationToken); @@ -128,21 +180,6 @@ public sealed class MediaProcessingBackgroundService( } } - /// Помечает ассет Processing и возвращает его расширение, либо null если обрабатывать нечего. - private async Task BeginProcessingAsync(Guid assetId, CancellationToken cancellationToken) - { - await using var scope = scopeFactory.CreateAsyncScope(); - var db = scope.ServiceProvider.GetRequiredService(); - - var asset = await db.MediaAssets.FirstOrDefaultAsync(x => x.Id == assetId, cancellationToken); - if (asset is null || asset.Status != MediaAssetStatus.Pending) - return null; - - asset.MarkProcessing(); - await db.SaveChangesAsync(cancellationToken); - return asset.OriginalExtension; - } - private async Task CompleteAsync( Guid assetId, MediaProcessingResult result,