1- using System ;
21using System . Collections . Concurrent ;
3- using System . Collections . Generic ;
42using System . Diagnostics ;
5- using System . IO ;
6- using System . Threading . Tasks ;
3+ using System . Runtime . ExceptionServices ;
74
8- using Microsoft . VisualBasic ;
5+ namespace ImageConvolution ;
96
10- namespace ImageConvolution
7+ public class ImageToProcess
118{
12- public class ImageToProcess
13- {
14- public string FilePath { get ; set ; } = "" ;
15- }
16-
17- public class ImageResult
18- {
19- public string OriginalFileName { get ; set ; } = "" ;
20- public double [ , ] ProcessedData { get ; set ; } = new double [ 0 , 0 ] ;
21- }
9+ public string FilePath { get ; set ; } = "" ;
10+ public double [ , ] ImageData { get ; set ; } = new double [ 0 , 0 ] ;
11+ }
2212
13+ public class ImageResult
14+ {
15+ public string OriginalFileName { get ; set ; } = "" ;
16+ public double [ , ] ProcessedData { get ; set ; } = new double [ 0 , 0 ] ;
17+ }
2318
24- public class AgentProcessor
19+ public class AgentProcessor
20+ {
21+ public static void ProcessImagesWithAgents ( string inputDir , string outputDir , int workerCount ,
22+ int queueCapacity = 4 , long memoryBudgetBytes = 0 )
2523 {
26- public static void ProcessImagesWithAgents ( string inputDir , string outputDir , int workerCount )
24+ ArgumentOutOfRangeException . ThrowIfNegativeOrZero ( workerCount ) ;
25+ ArgumentOutOfRangeException . ThrowIfNegativeOrZero ( queueCapacity ) ;
26+ ArgumentOutOfRangeException . ThrowIfNegative ( memoryBudgetBytes ) ;
27+ if ( ! Directory . Exists ( inputDir ) ) return ;
28+ ImageFiles . ValidateDirectories ( inputDir , outputDir ) ;
29+ Directory . CreateDirectory ( outputDir ) ;
30+ var stopwatch = Stopwatch . StartNew ( ) ;
31+ var files = ImageFiles . GetFiles ( inputDir ) ;
32+ using var slots = new SemaphoreSlim ( ImageFiles . GetInFlightLimit ( files , memoryBudgetBytes , 16 ) ) ;
33+ using var input = new BlockingCollection < ImageToProcess > ( queueCapacity ) ;
34+ using var output = new BlockingCollection < ImageResult > ( queueCapacity ) ;
35+ using var cancellation = new CancellationTokenSource ( ) ;
36+ var token = cancellation . Token ;
37+ Exception ? failure = null ;
38+ Task Start ( Action action ) => Task . Factory . StartNew ( ( ) =>
2739 {
28- Console . WriteLine ( $ "Запуск конвейера с { workerCount } агентами") ;
29- if ( ! Directory . Exists ( inputDir ) ) return ;
30- Directory . CreateDirectory ( outputDir ) ;
31-
32- var filesToProcessQueue = new BlockingCollection < ImageToProcess > ( 100 ) ;
33- var processedImagesQueue = new BlockingCollection < ImageResult > ( 100 ) ;
34-
35- var stopwatch = Stopwatch . StartNew ( ) ;
36- var readerAgent = Task . Run ( ( ) =>
40+ try { action ( ) ; }
41+ catch ( OperationCanceledException ) when ( token . IsCancellationRequested ) { }
42+ catch ( Exception ex )
3743 {
38-
39- foreach ( var file in Directory . GetFiles ( inputDir , "*jpg" ) )
44+ Interlocked . CompareExchange ( ref failure , ex , null ) ;
45+ cancellation . Cancel ( ) ;
46+ }
47+ } , CancellationToken . None , TaskCreationOptions . LongRunning , TaskScheduler . Default ) ;
48+ var reader = Start ( ( ) =>
49+ {
50+ try
51+ {
52+ foreach ( var file in files )
4053 {
41- filesToProcessQueue . Add ( new ImageToProcess { FilePath = file } ) ;
54+ slots . Wait ( token ) ;
55+ input . Add ( new ImageToProcess { FilePath = file , ImageData = ImageIO . LoadAsGrayscale ( file ) } , token ) ;
4256 }
43-
44- filesToProcessQueue . CompleteAdding ( ) ;
45- } ) ;
46-
47- var workerAgents = new List < Task > ( ) ;
48- for ( int i = 0 ; i < workerCount ; i ++ )
57+ }
58+ finally { input . CompleteAdding ( ) ; }
59+ } ) ;
60+ var workers = Enumerable . Range ( 0 , workerCount ) . Select ( _ => Start ( ( ) =>
61+ {
62+ foreach ( var item in input . GetConsumingEnumerable ( token ) )
4963 {
50- var worker = Task . Run ( ( ) =>
64+ var result = new ImageResult
5165 {
52- foreach ( var imageToProcess in filesToProcessQueue . GetConsumingEnumerable ( ) )
53- {
54- double [ , ] image = ImageIO . LoadAsGrayscale ( imageToProcess . FilePath ) ;
55- double [ , ] processedImage = ConvolutionProcessor . Convolve ( image , Kernels . BlurBox ) ;
56- var ResultForQue = new ImageResult
57- {
58- OriginalFileName = Path . GetFileName ( imageToProcess . FilePath ) ,
59- ProcessedData = processedImage
60- } ;
61- processedImagesQueue . Add ( ResultForQue ) ;
62- }
63- } ) ;
64- workerAgents . Add ( worker ) ;
66+ OriginalFileName = Path . GetFileName ( item . FilePath ) ,
67+ ProcessedData = ConvolutionProcessor . Convolve ( item . ImageData , Kernels . BlurBox )
68+ } ;
69+ item . ImageData = new double [ 0 , 0 ] ;
70+ output . Add ( result , token ) ;
6571 }
66-
67- var allWorkersTask = Task . WhenAll ( workerAgents ) . ContinueWith ( t =>
68- {
69- processedImagesQueue . CompleteAdding ( ) ;
70- } ) ;
71-
72-
73- var writerAgent = Task . Run ( ( ) =>
72+ } ) ) . ToArray ( ) ;
73+ var completion = Task . WhenAll ( workers ) . ContinueWith ( _ => output . CompleteAdding ( ) , TaskScheduler . Default ) ;
74+ var writer = Start ( ( ) =>
75+ {
76+ foreach ( var result in output . GetConsumingEnumerable ( token ) )
7477 {
75-
76- foreach ( var result in processedImagesQueue . GetConsumingEnumerable ( ) )
77- {
78- string savePath = Path . Combine ( outputDir , result . OriginalFileName ) ;
79- ImageIO . SaveImage ( result . ProcessedData , savePath ) ;
80- }
81- } ) ;
82-
83- Task . WaitAll ( readerAgent , allWorkersTask , writerAgent ) ;
84-
85- stopwatch . Stop ( ) ;
86- Console . WriteLine ( $ "\n Обработка конвейером завершена за { stopwatch . ElapsedMilliseconds } мс.") ;
87- }
78+ ImageIO . SaveImage ( result . ProcessedData , Path . Combine ( outputDir , result . OriginalFileName ) ) ;
79+ result . ProcessedData = new double [ 0 , 0 ] ;
80+ slots . Release ( ) ;
81+ }
82+ } ) ;
83+ Task . WaitAll ( reader , completion , writer ) ;
84+ if ( failure != null ) ExceptionDispatchInfo . Capture ( failure ) . Throw ( ) ;
85+ Console . WriteLine ( $ "Обработка конвейером ({ workerCount } агентов): { files . Length } файлов за { stopwatch . Elapsed . TotalMilliseconds : F1} мс.") ;
8886 }
89- }
87+ }
0 commit comments