Update media processing options: enhance .env.example and MediaOptions.cs to include MaxParallelTranscodes for improved transcoding performance. Refactor MediaProcessingBackgroundService to utilize parallel processing slots, ensuring efficient handling of multiple media assets simultaneously. Update comments for clarity on new configurations and processing logic.

This commit is contained in:
Leonid Pershin
2026-07-25 19:01:46 +03:00
parent 9cefc6a198
commit da0ac79079
3 changed files with 85 additions and 35 deletions
+8 -1
View File
@@ -44,8 +44,15 @@ Media__FfmpegPath=/usr/bin/ffmpeg
Media__FfprobePath=/usr/bin/ffprobe Media__FfprobePath=/usr/bin/ffprobe
# Максимальный размер загружаемого файла (20 ГБ). # Максимальный размер загружаемого файла (20 ГБ).
Media__MaxUploadBytes=21474836480 Media__MaxUploadBytes=21474836480
# Потоки ffmpeg — оставляем ядро API/раздаче (сервер 4 vCPU). # Потоки одного ffmpeg (0 — авто/все ядра).
Media__TranscodeThreads=3 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/ сканером, секунды. # Период опроса каталога inbox/ сканером, секунды.
Media__InboxScanSeconds=15 Media__InboxScanSeconds=15
# Нормализация громкости при обработке (EBU R128 loudnorm) — все ролики одинаковой громкости. # Нормализация громкости при обработке (EBU R128 loudnorm) — все ролики одинаковой громкости.
@@ -10,9 +10,15 @@ public sealed class MediaOptions
/// <summary>Максимальный размер загружаемого файла, байт (по умолчанию 20 ГБ).</summary> /// <summary>Максимальный размер загружаемого файла, байт (по умолчанию 20 ГБ).</summary>
public long MaxUploadBytes { get; init; } = 20L * 1024 * 1024 * 1024; public long MaxUploadBytes { get; init; } = 20L * 1024 * 1024 * 1024;
/// <summary>Число потоков ffmpeg — оставляем ядро API и раздаче (см. docs, 4 vCPU).</summary> /// <summary>Число потоков одного ffmpeg (0 — авто/все ядра). На многоядерных лучше сочетать с
/// <see cref="MaxParallelTranscodes"/>: несколько файлов по меньшему числу потоков дают больший
/// суммарный throughput, чем один файл на всех ядрах.</summary>
public int TranscodeThreads { get; init; } = 3; public int TranscodeThreads { get; init; } = 3;
/// <summary>Сколько файлов транскодировать одновременно (по умолчанию 1 — как было). Подбирается
/// под ядра: ориентир — <c>ядра / TranscodeThreads</c>, но не больше, чем позволяет память/диск.</summary>
public int MaxParallelTranscodes { get; init; } = 1;
/// <summary>Период опроса inbox/ сканером, секунды.</summary> /// <summary>Период опроса inbox/ сканером, секунды.</summary>
public int InboxScanSeconds { get; init; } = 15; public int InboxScanSeconds { get; init; } = 15;
@@ -1,3 +1,4 @@
using System.Collections.Concurrent;
using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Hosting;
@@ -9,11 +10,13 @@ using TeleWave.Domain.Media;
namespace TeleWave.Infrastructure.Media; namespace TeleWave.Infrastructure.Media;
/// <summary> /// <summary>
/// Единственный обработчик медиа: по одному ассету за раз прогоняет через ffmpeg. Источник истины — /// Обработчик медиа: прогоняет ассеты через ffmpeg, до <see cref="MediaOptions.MaxParallelTranscodes"/>
/// статус в БД: сервис в цикле берёт следующий Pending из базы и обрабатывает, а очередь лишь будит /// файлов одновременно. Источник истины — статус в БД: единственный диспетчер последовательно и
/// его без задержки. Поэтому рестарт/краш ничего не теряет — незавершённые задачи подхватываются из /// атомарно захватывает следующий Pending (помечает Processing), поэтому два транскода никогда не
/// БД (прерванные Processing на старте сбрасываются в Pending). БД-контекст держится короткими /// возьмут один ассет; сам транскод запускается в фоне с ограничением по числу слотов. Рестарт/краш
/// отрезками (пометить статус), сам транскод идёт вне scope, чтобы не держать соединение минутами. /// ничего не теряет — незавершённые подхватываются из БД (прерванные Processing на старте сбрасываются
/// в Pending). БД-контекст держится короткими отрезками (пометить статус), сам транскод идёт вне
/// scope, чтобы не держать соединение минутами.
/// </summary> /// </summary>
public sealed class MediaProcessingBackgroundService( public sealed class MediaProcessingBackgroundService(
IMediaProcessingQueue queue, IMediaProcessingQueue queue,
@@ -21,25 +24,59 @@ public sealed class MediaProcessingBackgroundService(
MediaPathResolver paths, MediaPathResolver paths,
IMediaProcessor processor, IMediaProcessor processor,
IOptions<StorageOptions> storageOptions, IOptions<StorageOptions> storageOptions,
IOptions<MediaOptions> mediaOptions,
ILogger<MediaProcessingBackgroundService> logger ILogger<MediaProcessingBackgroundService> logger
) : BackgroundService ) : BackgroundService
{ {
// Периодически перепроверяем БД, даже если сигнал не пришёл — страховка на любой случай. // Периодически перепроверяем БД, даже если сигнал не пришёл — страховка на любой случай.
private static readonly TimeSpan IdlePoll = TimeSpan.FromSeconds(30); private static readonly TimeSpan IdlePoll = TimeSpan.FromSeconds(30);
private readonly StorageOptions _storage = storageOptions.Value; private readonly StorageOptions _storage = storageOptions.Value;
private readonly int _maxParallel = Math.Max(1, mediaOptions.Value.MaxParallelTranscodes);
protected override async Task ExecuteAsync(CancellationToken stoppingToken) protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{ {
paths.EnsureDirectories(); paths.EnsureDirectories();
await ResetInterruptedAsync(stoppingToken); await ResetInterruptedAsync(stoppingToken);
// Слоты параллелизма: не запускаем больше _maxParallel транскодов одновременно.
using var slots = new SemaphoreSlim(_maxParallel, _maxParallel);
var inFlight = new ConcurrentDictionary<Task, byte>();
while (!stoppingToken.IsCancellationRequested) while (!stoppingToken.IsCancellationRequested)
{ {
try try
{ {
// Разобрать всю накопившуюся работу из БД. // Захватываем и раздаём по слотам всю накопившуюся работу из БД.
while (await NextPendingIdAsync(stoppingToken) is { } assetId) while (!stoppingToken.IsCancellationRequested)
await ProcessAsync(assetId, stoppingToken); {
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); using var wake = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
@@ -63,6 +100,16 @@ public sealed class MediaProcessingBackgroundService(
await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken); await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
} }
} }
// Даём уже запущенным транскодам корректно завершиться (или отмениться) на остановке.
try
{
await Task.WhenAll(inFlight.Keys.ToArray());
}
catch
{
// Ошибки/отмена отдельных задач уже залогированы внутри ProcessClaimedAsync.
}
} }
/// <summary>Сброс прерванных рестартом задач (Processing → Pending) на старте.</summary> /// <summary>Сброс прерванных рестартом задач (Processing → Pending) на старте.</summary>
@@ -82,26 +129,31 @@ public sealed class MediaProcessingBackgroundService(
await db.SaveChangesAsync(cancellationToken); await db.SaveChangesAsync(cancellationToken);
} }
/// <summary>Id самого раннего ассета в статусе Pending, либо null если работы нет.</summary> /// <summary>
private async Task<Guid?> NextPendingIdAsync(CancellationToken cancellationToken) /// Атомарно захватывает самый ранний Pending: помечает его Processing и возвращает (id, расширение),
/// либо null если работы нет. Вызывается только диспетчером последовательно, поэтому два транскода
/// не возьмут один ассет.
/// </summary>
private async Task<(Guid Id, string Extension)?> ClaimNextAsync(CancellationToken cancellationToken)
{ {
await using var scope = scopeFactory.CreateAsyncScope(); await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>(); var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
var id = await db.MediaAssets var asset = await db.MediaAssets
.Where(x => x.Status == MediaAssetStatus.Pending) .Where(x => x.Status == MediaAssetStatus.Pending)
.OrderBy(x => x.CreatedAt) .OrderBy(x => x.CreatedAt)
.Select(x => (Guid?)x.Id)
.FirstOrDefaultAsync(cancellationToken); .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) /// <summary>Обрабатывает уже захваченный (Processing) ассет: транскод → Ready/Failed.</summary>
private async Task ProcessClaimedAsync(Guid assetId, string extension, CancellationToken cancellationToken)
{ {
var extension = await BeginProcessingAsync(assetId, cancellationToken);
if (extension is null)
return;
try try
{ {
var result = await processor.ProcessAsync(assetId, extension, cancellationToken); var result = await processor.ProcessAsync(assetId, extension, cancellationToken);
@@ -128,21 +180,6 @@ public sealed class MediaProcessingBackgroundService(
} }
} }
/// <summary>Помечает ассет Processing и возвращает его расширение, либо null если обрабатывать нечего.</summary>
private async Task<string?> BeginProcessingAsync(Guid assetId, CancellationToken cancellationToken)
{
await using var scope = scopeFactory.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService<IAppDbContext>();
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( private async Task CompleteAsync(
Guid assetId, Guid assetId,
MediaProcessingResult result, MediaProcessingResult result,