View difference between Paste ID: V8ijvNeH and iPYWN1SZ
SHOW: | | - or go back to the newest paste.
1-
internal class ResourcePool : IDisposable
1+
using System;
2
using System.Linq;
3-
    private TimerCallback _callback;
3+
using System.Collections.Concurrent;
4-
    private bool _disposed;
4+
using System.Diagnostics;
5-
    private int _iDisposable;
5+
using System.Threading;
6-
    private TimeSpan _interval;
6+
using Core.Threading;
7-
    private int _max;
7+
8-
    private ArrayList _resources;
8+
namespace ConsoleApplication1
9-
    private Timer _timer;
9+
10
    class ResourcePoolExample
11-
    internal ResourcePool(TimeSpan interval, int max)
11+
12
13-
        this._interval = interval;
13+
        static void Main()
14-
        this._resources = new ArrayList(4);
14+
15-
        this._max = max;
15+
            
16-
        this._callback = new TimerCallback(this.TimerProc);
16+
            ExpensiveResource.InstanceCount = 0;
17
18
            int threadCount = 25;
19-
    [TargetedPatchingOptOut("Performance critical to inline this type of method across NGen image boundaries")]
19+
            int maxInstances = 2;
20-
    public void Dispose()
20+
21
            int min = threadCount- 2;
22-
        this.Dispose(true);
22+
            
23
            ThreadPool.SetMinThreads(min, min);
24
            //create factory to create ExpensiveResource instances
25-
    private void Dispose(bool disposing)
25+
            Func<ExpensiveResource> factory = () =>
26
            {
27-
        lock (this)
27+
                return new ExpensiveResource();
28
            };
29-
            if (!this._disposed)
29+
30
            ResourcePool<ExpensiveResource> pool = new ResourcePool<ExpensiveResource>(maxInstances, factory, 5000);
31-
                if (this._resources != null)
31+
32
            //event to signal when all threads are started
33-
                    foreach (IDisposable disposable in this._resources)
33+
            CountdownEvent allStarted = new CountdownEvent(threadCount);
34-
                    {
34+
35-
                        disposable.Dispose();
35+
            //event to signal when all threads are complete
36-
                    }
36+
            CountdownEvent allEnded = new CountdownEvent(threadCount);
37-
                    this._resources.Clear();
37+
38
            //event to signal all threads to go
39-
                if (this._timer != null)
39+
            ManualResetEventSlim goevent = new ManualResetEventSlim(false);
40
41-
                    this._timer.Dispose();
41+
            //Action that runs in every thread
42
            Action<int> testAction = (id) =>
43-
                this._disposed = true;
43+
44-
                if (disposing)
44+
45
                string indent = "".PadLeft(id, '-');
46-
                    GC.SuppressFinalize(this);
46+
                
47
                //signal this thread has started
48
                allStarted.Signal();
49
                Console.Write(".");
50
                //wait for every thread
51
                goevent.Wait();
52-
    ~ResourcePool()
52+
53
                //delay each thread to simulate reality.
54-
        this.Dispose(false);
54+
                Random rnd = new Random(Guid.NewGuid().ToString().GetHashCode());
55
                Thread.Sleep(rnd.Next(10, 100));
56
57-
    internal object RetrieveResource()
57+
                ExpensiveResource resource = null;
58
59-
        object obj2 = null;
59+
60-
        if (this._resources.Count != 0)
60+
61
62-
            lock (this)
62+
                    //we grab our resource
63
                    resource = pool.Take();
64-
                if (this._disposed)
64+
65
                    Console.WriteLine(indent + "Action{0} Take -> Instance_{1}", id, resource.InstanceId);
66-
                    return obj2;
66+
67
                    //execute process
68-
                if (this._resources.Count == 0)
68+
                    resource.ExpensiveProcess();
69
70-
                    return null;
70+
71
                catch (Exception error)
72-
                obj2 = this._resources[this._resources.Count - 1];
72+
73-
                this._resources.RemoveAt(this._resources.Count - 1);
73+
                    Console.WriteLine("ERROR ExpensiveProcess: {0}", error.Message);
74-
                if (this._resources.Count < this._iDisposable)
74+
75
                finally
76-
                    this._iDisposable = this._resources.Count;
76+
77
                    //return our resource to the pool when finished
78
                    pool.Return(resource);
79
                    Console.WriteLine(indent + "Action{0} Return -> Instance_{1}", id, resource.InstanceId);
80-
        return obj2;
80+
81
                    //signal when finished
82
                    allEnded.Signal();
83-
    internal void StoreResource(IDisposable o)
83+
84
            };
85-
        lock (this)
85+
86
            Console.WriteLine("Starting {0:N0} threads", threadCount);
87-
            if (!this._disposed && (this._resources.Count < this._max))
87+
88
            //create our worker threads
89-
                this._resources.Add(o);
89+
            for (int i = 0; i < threadCount; i++)
90-
                o = null;
90+
91-
                if (this._timer == null)
91+
                int workerId = i + 1;
92
93-
                    this._timer = new Timer(this._callback, null, this._interval, this._interval);
93+
                ThreadPool.QueueUserWorkItem((state) =>
94
                {
95
                    testAction((int)state);
96
                }, workerId);
97-
        if (o != null)
97+
98
                
99-
            o.Dispose();
99+
100
            Console.WriteLine();
101
102
            //wait for all threads  to startup
103-
    private void TimerProc(object userData)
103+
            allStarted.Wait();
104
            Console.WriteLine();
105-
        IDisposable[] array = null;
105+
            Console.WriteLine("All {0:N0} threads have started",threadCount);
106-
        lock (this)
106+
            //fire off all threads
107
            goevent.Set();
108-
            if (!this._disposed)
108+
109
            //wait until all threads are completed
110-
                if (this._resources.Count == 0)
110+
            allEnded.Wait();
111
112-
                    if (this._timer != null)
112+
            //uncomment to use 'Parallel lib'
113-
                    {
113+
            //Parallel.For(0,threadCount,testAction);
114-
                        this._timer.Dispose();
114+
115-
                        this._timer = null;
115+
            Console.WriteLine();
116-
                    }
116+
            Console.WriteLine("threadCount {0}", threadCount);
117-
                    return;
117+
            Console.WriteLine("pool.MaxInstances: {0}", pool.MaxInstances);
118
            Console.WriteLine("pool.InstancesCreated: {0}", pool.InstancesCreated);
119-
                array = new IDisposable[this._iDisposable];
119+
            Console.WriteLine("ExpensiveResource.InstanceCount: {0}", ExpensiveResource.InstanceCount);
120-
                this._resources.CopyTo(0, array, 0, this._iDisposable);
120+
            Console.WriteLine("pool.Count: {0}", pool.Count);
121-
                this._resources.RemoveRange(0, this._iDisposable);
121+
            Console.WriteLine();
122-
                this._iDisposable = this._resources.Count;
122+
            Console.WriteLine("Test Finished");
123
        }
124
    }
125-
        if (array != null)
125+
    class ExpensiveResource
126
    {
127-
            for (int i = 0; i < array.Length; i++)
127+
128
        //global instances created
129
        public static int InstanceCount;
130
131-
                    array[i].Dispose();
131+
        //resource id
132
        public int InstanceId { get; private set; }
133-
                catch
133+
134
        public ExpensiveResource()
135
        {
136
137
            InstanceId = Interlocked.Increment(ref InstanceCount);
138
            Console.WriteLine("Created Instance #: {0}", InstanceId);
139
        }
140
141
        public void ExpensiveProcess()
142
        {
143
            //simulate long process
144
            Random rnd = new Random(Guid.NewGuid().ToString().GetHashCode());
145
            int st = rnd.Next(500, 1000);
146
            Thread.Sleep(st);
147
        }
148
    }
149
}
150
151
152
namespace Core.Threading {
153
154
    /// <summary>
155
    /// Resource pool that uses a <see cref="System.Collections.Concurrent.BlockingCollection<T>"/> that limits the use of specified resource.
156
    /// If an item in the pool is requested and the instance count is under the maximum a new instance of the resource is created via <see cref="Func<T>"/> parameter.
157
    /// </summary>
158
    /// <example>
159
    /*
160
         class ResourcePoolExample {
161
         
162
        static void Main() {
163
164
            ExpensiveResource.InstanceCount = 0;
165
166
            int threadCount = 8;
167
            int maxInstances = 2;
168
169
            //create factory to create ExpensiveResource instances
170
            Func<ExpensiveResource> factory = () => {
171
                return new ExpensiveResource();
172
            };
173
174
            ResourcePool<ExpensiveResource> pool = new ResourcePool<ExpensiveResource>(maxInstances,factory,5000);
175
            
176
            //event to signal when all threads are complete
177
            CountdownEvent cde = new CountdownEvent(threadCount);
178
179
            //event to signal all threads to go
180
            ManualResetEventSlim goevent = new ManualResetEventSlim(false);
181
182
            //Action that runs in every thread
183
            Action<int> testAction = (id) => {
184
185
                string indent = "".PadLeft(id,'-');
186
                
187
                //wait for every thread
188
                goevent.Wait();
189
190
                //delay each thread to simulate reality.
191
                Random rnd = new Random(Guid.NewGuid().ToString().GetHashCode());
192
                Thread.Sleep(rnd.Next(10,100));
193
194
                ExpensiveResource resource = null;
195
196
                try {
197
198
                    //we grab our resource
199
                    resource = pool.Take();
200
201
                    Console.WriteLine(indent + "Action{0} Take -> Instance_{1}",id,resource.InstanceId);
202
203
                    //execute process
204
                    resource.ExpensiveProcess();
205
206
                }
207
                catch(Exception error) {
208
                    Console.WriteLine("ERROR ExpensiveProcess: {0}",error.Message);
209
                }
210
                finally {
211
                    //return our resource to the pool when finished
212
                    pool.Return(resource);
213
                    Console.WriteLine(indent + "Action{0} Return -> Instance_{1}",id,resource.InstanceId);
214
215
                    //signal when finished
216
                    cde.Signal();
217
                }
218
            };
219
220
            //create our worker threads
221
            for(int i = 0;i < threadCount;i++) {
222
                int workerId = i + 1;
223
224
                ThreadPool.QueueUserWorkItem((state) => {
225
                    testAction((int)state);
226
                },workerId);
227
            }
228
229
            Thread.Sleep(10);
230
     
231
            //fire off all threads
232
            goevent.Set();
233
234
            //wait until all threads are completed
235
            cde.Wait();
236
237
            //uncomment to use 'Parallel lib'
238
            //Parallel.For(0,threadCount,testAction);
239
240
            Console.WriteLine();
241
            Console.WriteLine("threadCount {0}",threadCount);
242
            Console.WriteLine("pool.MaxInstances: {0}",pool.MaxInstances);
243
            Console.WriteLine("pool.InstancesCreated: {0}",pool.InstancesCreated);
244
            Console.WriteLine("ExpensiveResource.InstanceCount: {0}",ExpensiveResource.InstanceCount);
245
            Console.WriteLine("pool.Count: {0}",pool.Count);
246
            Console.WriteLine();
247
            Console.WriteLine("Test Finished");
248
        }
249
    }
250
    class ExpensiveResource {
251
252
        //global instances created
253
        public static int InstanceCount;
254
255
        //resource id
256
        public int InstanceId { get; private set; }
257
258
        public ExpensiveResource() {
259
            
260
            InstanceId = Interlocked.Increment(ref InstanceCount);
261
            Console.WriteLine("Created Instance #: {0}",InstanceId);
262
        }
263
264
        public void ExpensiveProcess() {
265
            //simulate long process
266
            Random rnd = new Random(Guid.NewGuid().ToString().GetHashCode());
267
            int st = rnd.Next(500,1000);
268
            Thread.Sleep(st);
269
        }
270
    }
271
     */
272
    /// </example>
273
    /// <typeparam name="T">class</typeparam>
274
    public class ResourcePool<T> where T:class {
275
276
        protected readonly object _syncroot = new object();
277
        protected readonly BlockingCollection<T> _items;
278
        protected readonly Func<T> _factory;
279
        readonly int _timeout;
280
        volatile int _maxInstances;
281
        volatile int _created;
282
283
        /// <summary>Instances that were created</summary>
284
        public int InstancesCreated { get { return _created; } }
285
286
        /// <summary>Maximum instances that could be created</summary>
287
        public int MaxInstances { get { return _maxInstances; } }
288
289
        /// <summary>Gets the number of current objects available in the pool</summary>
290
        public int Count { get { return _items.Count; } }
291
292
        ///<summary>Gets the wait time if all resources are taken and no more can be created, default is -1 (infinity)</summary>
293
        public int Timeout { get { return _timeout; } }
294
295
        /// <summary>
296
        /// ctr
297
        /// </summary>
298
        /// <param name="maxInstances">Maximum number of instances that are to be created</param>
299
        /// <param name="factory">Factory function to create an instance of the specified Type </param>
300
        /// <param name="timeout">The wait time if all resources are taken and no more can be created, default is -1 (infinity)</param>
301
        public ResourcePool(int maxInstances,Func<T> factory,int timeout = -1) {
302
303
            if(factory == null)
304
                throw new ArgumentNullException("factory");
305
306
            this._maxInstances = Math.Max(1,maxInstances);
307
            this._items = new BlockingCollection<T>(_maxInstances);
308
            this._factory = factory;
309
            this._timeout = timeout;
310
        }
311
312
        /// <summary>Increments the MaxInstances value</summary>
313
        public void IncrementMaxInstances(int count) {
314
            _maxInstances += count;
315
        }
316
317
        protected virtual void OnCreated(T instance) {
318
            return;
319
        }
320
        protected virtual bool CanReturn(T instance) {
321
            return true;
322
        }
323
324
        /// <summary>
325
        /// Returns an object back to the resource pool to be available the next time <see cref="Take"/> is called 
326
        /// </summary>
327
        /// <param name="item">The object to be placed back into the pool</param>
328
        /// <returns>True if object is added back to the pool, otherwise false</returns>
329
        public virtual bool Return(T item) {
330
            if(item == null)
331
                return false;
332
333
            bool added = false;
334
            if(CanReturn(item)) {
335
                added = _items.TryAdd(item);
336
            }
337
#if DEBUG
338
            Debug.WriteLine("Return-> hash: {0} - added: {1}",item.GetHashCode(),added);
339
#endif
340
            return added;
341
        }
342
343
        /// <summary>
344
        /// Either takes an existing object from the pool, or it creates a new instance if the <see cref="InstancesCreated"/> have not exceed the <see cref="MaxInstances"/> value.
345
        /// </summary>
346
        /// <returns>An existing object from the pool or a new instnace</returns>
347
        public virtual T Take() {
348
349
#if DEBUG
350
            int hash = -1;
351
            bool added = false;
352
             
353
#endif
354
            bool created = false;
355
            T item = null;
356
357
            //take out the resource immediately
358
            if(_items.TryTake(out item,0))
359
                return item;
360
361
            lock(_syncroot) {
362
363
                //create resource if we can
364
                if(_created < _maxInstances) {
365
366
                    item = _factory();
367
368
                    _created++;
369
                    _items.Add(item);
370
                    created = true;
371
                }
372
            }
373
374
            if(created) {
375
                OnCreated(item);
376
                return item;
377
            }
378
379
            if(_created >= _maxInstances) {
380
381
                //all resources are created, wait until there is one available
382
                if(!_items.TryTake(out item,_timeout))
383
                    //we timed out
384
                    throw new TimeoutException();
385
            }
386
387
#if DEBUG
388
            if(item != null)
389
                hash = item.GetHashCode();
390
             
391
#endif
392
            return item;
393
        }
394
    }
395
}