[Java] Hoe weten dat threads klaar zijn

Pagina: 1
Acties:
  • 150 views sinds 30-01-2008
  • Reageer

  • Glimi
  • Registratie: Augustus 2000
  • Niet online

Glimi

Designer Drugs

Topicstarter
(overleden)
Heren,

Ik ben op het moment bezig met het maken van een quicksort, op het moment alleen nog voor integers maar meer dan C&P zou het niet moeten zijn voor float/doubles. De quicksort is het probleem niet (hij loopt zelfs wat sneller dan Arrays.sort(int[]) en werkt perfect, maar nu moet er eigenlijk ook een threaded variant van de quicksort komen.

Daar zitten de problemen. Ik kom er namelijk niet uit hoe ik er voor kan zorgen dat de gehele array gesorteerd is, voordat ik retourneer.
Voordat ik daar meer over kan vertellen eerst wat code.



Eerst de threadpool waar de ThreadedQuicksort gebruik van maakt. Met deze klasse tracht ik het maximaal aantal threads te beheersen en enige administratie mogelijk te maken. Tevens voorkom je natuurlijk de overhead van meerdere creëaties door Threads te hergebruiken
Java:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
import java.util.*;

/**
 * Created by IntelliJ IDEA.
 * User: Glimi Development
 * Date: Sep 22, 2003
 * Time: 6:45:59 PM
 * To change this template use Options | File Templates.
 */

public class Threadpool {

    private int _maxSize;
    private int _currentActive = 0;
    private List _waiting = Collections.synchronizedList( new ArrayList() );

    public Threadpool( final int p_maxSize ) {

        initPool( p_maxSize );
        setMaxSize( p_maxSize );
    }

    synchronized
    public void setMaxSize( final int p_maxSize ) {

        if ( p_maxSize <= 0 )
            throw new IllegalArgumentException( "The maxsize cannot be negative : " + p_maxSize );

        _maxSize = p_maxSize;
    }

    public int getMaxSize() {

        return _maxSize;
    }

    public int getCurrentActive() {

        return _currentActive;
    }

    private void increaseCurrentActive() {

        ++_currentActive;
    }

    private void decreaseCurrentActive() {

        --_currentActive;
    }

    synchronized
    public void runJob( final Runnable p_job ) {

        final PooledThread l_thread;

        while ( getMaxSize() <= getCurrentActive() )
            try {

                wait();
            } catch ( InterruptedException ie ) {

                ie.printStackTrace();
            }

        if ( _waiting.isEmpty() ) {

            l_thread = new PooledThread();
            l_thread.start();
        } else
            l_thread = (PooledThread) _waiting.get( _waiting.size() - 1 );

        l_thread.wakeAndRunJob( p_job );
        increaseCurrentActive();
    }

    synchronized
    private boolean push( final PooledThread p_finishedThread ) {

        boolean l_letTreadLive = true;

        if ( getMaxSize() < _waiting.size() + getCurrentActive() )
            l_letTreadLive = false;

        if ( l_letTreadLive ) {
            _waiting.add( _waiting.size(), p_finishedThread );
            notify();
        }

        decreaseCurrentActive();
        return l_letTreadLive;
    }

    private void initPool( final int p_initialSize ) {

        for ( int i = 0; i < p_initialSize; ++i )
            _waiting.add( i, new PooledThread() );
    }

    class PooledThread extends Thread {

        private Runnable _job = null;

        public PooledThread() {

            this( null );
        }

        public PooledThread( final Runnable p_job ) {

            setJob( p_job );
        }

        private void setJob( final Runnable p_job ) {

            _job = p_job;
        }

        synchronized
        public void wakeAndRunJob( final Runnable p_job ) {

            setJob( p_job );
            notify();
        }

        synchronized
        public void run() {

            boolean l_stop = false;

            while ( !l_stop ) {

                // If there is no job to do, go to sleep until notifyed by wakeAndRunJob
                //
                if ( _job == null )
                    try {

                        wait();
                    } catch ( InterruptedException ie ) {

                        ie.printStackTrace();
                        continue;
                    }

                // If the thread woke up and has data
                //
                if ( _job != null )
                    _job.run();

                _job = null;
                l_stop = push( this );
            }
        }
    }
}
Her komt er kort op neer dat als er nog threads vrij zijn, deze de gegeven job in de maag gesplitst krijgen en bij het voltooien van die job, zich weer aanmelden bij de pool.
Als er geen threads vrij zijn en het maxAantalThreads is bereikt, dan wordt de thread op hold gezet en weer wakker gemaakt als er een thread zich weer aanmeld.

ThreadedQuicksort
De threadedquicksort sorteert de gegeven array van ints dmv het aanmaken van een InternalQuickSortJob voor elke range die gesorteerd moet worden. Echter, na het aanmaken van zo'n job, returned hij wel gelijk. Zodra quicksort() dus returned is de boel nog niet gesorteerd.
Java:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
public class ThreadedQuickSort extends QuickSort {

    private final static int DEFAULT_MAX_SIZE = 10;
    private Threadpool _pool;

    public ThreadedQuickSort( ) {

        this( new Threadpool( DEFAULT_MAX_SIZE ) );
    }

    public ThreadedQuickSort( final Threadpool p_pool ) {

        setPool( p_pool );
    }

    public void setPool( final Threadpool p_pool ) {

        if( p_pool == null )
            throw new NullPointerException( "A null Threadpool cannot be used" );

        _pool = p_pool;
    }

    public Threadpool getPool( final int p_maxSize ) {

        return _pool;
    }

    public void quicksort( final int[] p_source, final int p_offset, final int p_length ) {

        super.quicksort( p_source, p_offset, p_length );
    }

    protected void recurse( final int[] p_source, final int p_offset, final int p_length ) {

        if( p_length > 1) {

            InternalQuickSortJob l_job = new InternalQuickSortJob( p_source, p_offset, p_length);
            _pool.runJob( l_job );
        }
    }

    private class InternalQuickSortJob implements Runnable {

        private int[]             _source;
        private int               _offset;
        private int               _length;

        public InternalQuickSortJob( final int[] p_source, final int p_offset, final int p_length ) {

            _source = p_source;
            _offset = p_offset;
            _length = p_length;
        }

        public void run(){

            quicksort( _source, _offset, _length);
        }
    }
}

Hoe kan ik er nou voor zorgen dat bij de eerste aanroep van quicksort() de boel pas retourneerd als de array gesorteerd is?

Geprobeerd
Het volgende is geprobeerd:
• Dispatch een event vanaf de Threadpool als de boel leeg is. Echter doordat het event wordt afgehandeld in een andere thread krijg je problemen met notify() en wait()
• Gebruik join() heeft geen zin, de threads zijn volledig afgeschermd en dat moet eigenlijk zo blijven
• Ik heb bij het aanroepen van quicksort() een Semaphore op 1- p_length geinitaliseerd. Elke keer als recurse() een length van 1 of kleiner tegenkomt, geeft hij een release() op de Semaphore. quicksort() doet dan na de aanroep van een parition (welke feitelijk super.quicksort() aanroept) een aquire() waardoor hij pas kan returnen nadat alle p_length stukken een release() hebben gedaan.
Echter dan kom ik in een deathlock :/

Heeft iemand nog suggesties voor me? Moet de structuur om, of zie ik iets over het hoofd? Alvast bedankt voor het lezen iig :)

Verwijderd

Ah, Glimi heeft een threadpool gemaakt. Daar dacht ik ook direct aan toen ik de opdracht hoorde. De meeste mensen hebben echter een eenvoudigere oplossing gebruikt ;-)

Het is voor de opdracht niet nodig om op terminatie te wachten, maar is natuurlijk wel mooier. Ik weet zo snel wel een manier, maar die is niet zo mooi: je maakt een aparte thread die periodiek de grootte van de pool checkt en het zaakje vrijgeeft als de pool leeg is (ik ga er vanuit dat er minstens een thread in de pool is zolang het sorteren loopt). Helaas is dit ook nog eens applicatie specifiek, maar daar kan je met het observer pattern wel wat aan doen denk ik.

Ik heb er verder niet veel aandacht aan geschonken. Als je wilt kan je maandag naar het practicum komen en kunnen we er even over babbelen, maar dat is wel wat kort dag voor je.

[ Voor 15% gewijzigd door Verwijderd op 04-10-2003 13:32 . Reden: aanvulling ]


  • Infinitive
  • Registratie: Maart 2001
  • Laatst online: 10-08 15:15
Ik weet nog wel een oplossing. Misschien niet de mooiste, maar dat laat ik aan jezelf om te beoordelen.

Je wilt weten wanneer je in jobjes verdeelde algoritme afgelopen is. We stellen alleen dat je algoritme klaar is als en alleen als er geen actieve jobs zijn en ook geen jobs zitten te wachten. Zolang je algoritme runt moet deze eis gewaarborgt blijven, maar dat is natuurlijk vrij makkelijk in te zien. Merk op dat dit een situatie is die gecontroleerd kan worden nadat een jobje klaar is.

Verder nemen we een extra thread, de hoofdthread. Deze thread doet alles maar toch vrij weinig: het zet een eerst semaphore op scherp en queued je hoofdjob. Vervolgens gaat deze thread wachten totdat de semaphore vrijgegeven wordt. Tot slot geef het op een of andere manier (afhankelijk van je implementatie) aan elke thread van de pool door dat zij kunnen termineren en wacht totdat deze dood zijn via de join() methode. Wanneer je de semaphore vrijgeeft op het moment dat de thread die het laatste jobje zojuist heeft uitgevoerd echter komt dat het algorime afgelopen heeft, zal je hoofdthread het hele zaakje netjes laten termineren. Op zich hoeft de hoofdthread geen extra thread te zijn, maar een ding is wel belangrijk: er mogen geen jobjes op uitgevoerd worden.

Aangezien java zelf geen ingebakken semaphoren heeft kan je deze maken via een monitor. Verder kan je met een monitor je threadpool maken. Daarvoor zul je dus minstens twee verschillende monitoren gebruiken. Merk op dat je met een 'single producer, multiple consumer buffer met stop functie' (queue) je threadpool ook wat eenvoudiger kan implementeren.

putStr $ map (x -> chr $ round $ 21/2 * x^3 - 92 * x^2 + 503/2 * x - 105) [1..4]


  • Glimi
  • Registratie: Augustus 2000
  • Niet online

Glimi

Designer Drugs

Topicstarter
(overleden)
Verwijderd schreef op 04 October 2003 @ 13:22:
Ik heb er verder niet veel aandacht aan geschonken. Als je wilt kan je maandag naar het practicum komen en kunnen we er even over babbelen, maar dat is wel wat kort dag voor je.
Ik zal in ieder geval morgenochtend even de rest van de opdracht even doen. Ik kom morgen dan wel naar practicum om je lastig te vallen. Dat lukt allemaal wel.
Infinitive schreef op 04 October 2003 @ 23:39:
Ik weet nog wel een oplossing. Misschien niet de mooiste, maar dat laat ik aan jezelf om te beoordelen.
Nou veel andere oplossingen heb ik niet gezien, dus veel competitie heeft het niet :+
Algoritme, welke ik niet helemaal volg
Zoals je ziet ben ik je een beetje kwijt. Misschien dat we er morgen nog even over kunnen babbelen. Maar wat ik er van begrijp wil je een Semafoor laten aanmaken in de eerste aanroep, voordat enige recurrente aanroep heeft plaats gevonden, en dat bij de laatste recurrente aanroep de semafoor pas 'up/beschikbaar' wordt gezet.
Daarbij speelt wel het probleem dat je niet weet waarneer het klaar is.

Dit heb ik (ongeveer) geprobeerd met het volgende:
• Move de inhoud van ThreadedQuickSort.quicksort() naar een methode genaamd partition
• Laat in ThreadedQuickSort.quicksort() een semafoor initialiseren op 1-n, waarbij n gelijk is aan de lengte van het te sorteren stuk. Immers dan kan na n maal een release() pas een aquire() lukken
• Laat ThreadedQuickSort.quicksort() partition() aanroepen.
• Laat ThreadedQuickSort.InternalQuickSortJob.run() ook partition() aanroepen ipv quicksort()
• Laat recurse() een Semafoor.release() doen als hij een stuk van lengte <= 1 moet gaan recursen.

Echter toen kwam ik in een (voor mij) overklaarbare deathlock terecht :(
Aangezien java zelf geen ingebakken semaphoren heeft kan je deze maken via een monitor. Verder kan je met een monitor je threadpool maken. Daarvoor zul je dus minstens twee verschillende monitoren gebruiken. Merk op dat je met een 'single producer, multiple consumer buffer met stop functie' (queue) je threadpool ook wat eenvoudiger kan implementeren.
Ik heb de concurrency library van Doug Lea gebruikt. Deze komt in Java 1.5 (enigzins aangepast) dus lijkt me een juiste keuze qua toekomstzicht. Deze bevat ook een class Semafoor. Zie http://gee.cs.oswego.edu/...til/concurrent/intro.html

[ Voor 3% gewijzigd door Glimi op 05-10-2003 22:33 ]


  • Soultaker
  • Registratie: September 2000
  • Laatst online: 20-08 00:10
Als je InternalQuickSortJobs's nu wisten bij welke parent Job ze hoorden, zou het probleem snel opgelost zijn. Het probleem zit 'm er volgens mij in dat zodra recurse wordt aangeroepen (door de originele implementatie van quicksort in de QuickSort-klasse, als ik het goed begrijp) je geen idee meer hebt voor welke taak je een subthread gaat spawnen, noch retourneer je informatie over de gespawnde subthread aan de caller.

Het lijkt me dus dat je de koppeling tussen taak en subtaak nooit kunt maken, tenzij je gebruik maakt van 'globale' informatie zoals de identificatie van de huidige thread (geen idee of en hoe die beschikbaar is in Java, trouwens). Ik zou me kunnen voorstellen dat je thread pool onthoudt welke thread een job aanbiedt en op basis van die gegevens een 'waitForChildren'-achtige methode aan zou kunnen bieden. Dat gaat natuurlijk wel ten koste van een stukje efficientie in de thread pool implementatie (je moet iets meer gegevens bijhouden), maar het is toch al Java dus waar hebben we het over. ;)

[ Voor 3% gewijzigd door Soultaker op 05-10-2003 23:02 ]


Verwijderd

hmm ik weet niet zoveel meer van java (wordt weinig gebruikt bij mij op het werk), maar als je nou es alle jobs/subjobs in een queue gooide? Dan hoef je alleen maar te checken of de queue leeg is (hij is dan klaar). Wel pas een job uit de queue gooien als hij klaar is natuurlijk...

Heb je gelijk een stukje administratie over alle lopende jobs/threads.

edit : okee eventjes ietsje verduidelijken...

Je past je threadpool aan zodat je er een administratie van je jobs in bijhoud (in een vector/queue, weet niet wat java daarvoor biedt). Die pool heeft al een main thread. Deze laat je een job uit zn queue uitvoeren in een van de threads in de pool. De recurse methode laat je een job toevoegen aan de jobqueue van de threadpool. In tegenstelling tot mijn eerdere statement wel gelijk uit de queue verwijderen (anders ga je dingen dubbel uitvoeren). Als er nu geen jobs in de queue meer zijn en alle threads in de pool niks meer aan het doen zijn is de hoofdjob klaar. Om dat naar je main thread te signalen kun je een event gebruiken (ik denk in C++ termen, dus een synchronisatie event wat je kunt setten en resetten).

[ Voor 53% gewijzigd door Verwijderd op 05-10-2003 23:26 ]


  • Infinitive
  • Registratie: Maart 2001
  • Laatst online: 10-08 15:15
Ik heb het niet helemaal goed doorgelezen, maar volgens mij bedoelde ik hetzelfde als Akhorahil hierboven schreef.

Zo'n event kan je eenvoudig met een monitor maken. Daar heb je alleen synchronized, wait&notify en een boolean voor nodig.

Glimi, als je externe biblotheken gebruikt of een speciale/erg nieuwe jdk versie nodig hebt, dat je dit wel even documenteerd? Nakijkers vinden het bijv. niet erg leuk als er direct een of andere nullpointer exception om je oren vliegt als er een bibliotheek niet geinstalleerd is ;)

putStr $ map (x -> chr $ round $ 21/2 * x^3 - 92 * x^2 + 503/2 * x - 105) [1..4]


  • Glimi
  • Registratie: Augustus 2000
  • Niet online

Glimi

Designer Drugs

Topicstarter
(overleden)
Soultaker schreef op 05 October 2003 @ 23:02:
Als je InternalQuickSortJobs's nu wisten bij welke parent Job ze hoorden, zou het probleem snel opgelost zijn. Het probleem zit 'm er volgens mij in dat zodra recurse wordt aangeroepen (door de originele implementatie van quicksort in de QuickSort-klasse, als ik het goed begrijp) je geen idee meer hebt voor welke taak je een subthread gaat spawnen, noch retourneer je informatie over de gespawnde subthread aan de caller.
Inderdaad, er is een bepaalde staat van statelessness waardoor al die jobs niet weten of ze überhaupt een parent hebben die ik opzich wel mooi vind.
Het lijkt me dus dat je de koppeling tussen taak en subtaak nooit kunt maken, tenzij je gebruik maakt van 'globale' informatie zoals de identificatie van de huidige thread (geen idee of en hoe die beschikbaar is in Java, trouwens).
Ja. Je kunt threads identificeren aan namen, Threadgroups ed. ed. :) Probleem in mijn implementatie is echter wel dat Threads hergebruikt worden voor de jobs en dat de jobs opzich niet blocken, wat dus een beindiging van de Job kan betekenen voordat een andere Job kan communiceren ermee. Echter dat zou opzich te verhelpen zijn.
Ik zou me kunnen voorstellen dat je thread pool onthoudt welke thread een job aanbiedt en op basis van die gegevens een 'waitForChildren'-achtige methode aan zou kunnen bieden. Dat gaat natuurlijk wel ten koste van een stukje efficientie in de thread pool implementatie (je moet iets meer gegevens bijhouden), maar het is toch al Java dus waar hebben we het over. ;)
Als je voor elke iteratie 2 threads gaat jongen en daar op gaat wachten (ok, dit is te besparen tot één maar niet in deze implementatie) dan kom je toch uit op nlogn threads op het eind?

Stel Job X jongt Job Y en roept een eventuele waitForChildren() aan. X kan dus pas voltooien als Job Y voltooid is, maar die kan pas voltooien als zijn gejongde Thread Z voltooid is. Kortom elke job levert minimaal één extra job op waarop gewacht moet worden en opzich krijg je gewoon een sequentïele handel waarbij je niet de threads kan limieteren. Of zie ik het fout?


Verwijderd schreef op 05 October 2003 @ 23:13:
edit : okee eventjes ietsje verduidelijken...

Je past je threadpool aan zodat je er een administratie van je jobs in bijhoud (in een vector/queue, weet niet wat java daarvoor biedt). Die pool heeft al een main thread. Deze laat je een job uit zn queue uitvoeren in een van de threads in de pool. De recurse methode laat je een job toevoegen aan de jobqueue van de threadpool. In tegenstelling tot mijn eerdere statement wel gelijk uit de queue verwijderen (anders ga je dingen dubbel uitvoeren). Als er nu geen jobs in de queue meer zijn en alle threads in de pool niks meer aan het doen zijn is de hoofdjob klaar. Om dat naar je main thread te signalen kun je een event gebruiken (ik denk in C++ termen, dus een synchronisatie event wat je kunt setten en resetten).
Zo is de implementatie dus ongeveer ook gewoorden. Gekozen voor een queue om FIFO te kunnen garanderen en gekozen om de boel via een normaal event te laten handelen (die dus in z'n eigen thread loopt) en threadsafe te maken door deze een Semafoor up te gooien.

Misschien is de code wat duidelijker dan dit vage gezang.

ThreadPool
Werkt niet met een interne thread voor het uitvoeren van Jobs, maar gaat jobs processen zodra ze binnen komen en/of een thread zich aanmeld door middel van push()
Merk trouwens op dat een van mijn eerste problemen in dit topic zichzelf opeens oploste toen ik de door de pool aangemaakte threads ook werkelijk startte 8)7
Java:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
public class FIFOThreadPool extends AbstractThreadPool { 

    // The maximum amount of active threads that can run simoultaniously 
    // 
    private int _maxSize; 

    // Variable used to count the current running jobs 
    // 
    private int _currentActive = 0; 

    // A Stack used to store the Threads that are waiting for jobs 
    // A Stack is used purely for easy access 
    // 
    private final Stack _inactiveThreads = new Stack(); 

    // A queue to store the Jobs that are waiting for scheduling 
    // A queue is used to provide the FIFO scheduling garuantee 
    // 
    private final LinkedList _waitingJobs = new LinkedList(); 

    /** 
     * Create a new <code>FIFOThreadPool</code> with the given maximum 
     * of active jobs 
     * @param p_maxSize the maximum of active jobs 
     */ 
    public FIFOThreadPool( final int p_maxSize ) { 

        initPool( p_maxSize ); 
        setMaxSize( p_maxSize ); 
    } 

    /** 
     * Set the maximum of active jobs that this pool can provide. 
     * <p> 
     * This method is synchronized to make sure that calls made to this 
     * function, are finished befor a new job 
     * is scheduled and/or another thread can finish it. 
     * @param p_maxSize maximum active jobs 
     */ 
    synchronized 
    public void setMaxSize( final int p_maxSize ) { 

        _maxSize = p_maxSize; 

        /* If there are jobs waiting, but there is no activity 
         * (probably because the _maxSize was set to 0 before this call) 
         * then pick the last job form the waitingJobs and schedule it. 
         * 
         * It is guaranteed that the job taken from the queue is set back 
         * to it's place or scheduled (if it was the only job!) 
         * because this method and the runJob() 
         * method share the same monitor. 
         */ 
        if ( !_waitingJobs.isEmpty() && _currentActive == 0 ) 
            runJob( (Runnable) _waitingJobs.removeLast() ); 
    } 

    /** 
     * Return the maximum of active jobs that this pool can provide. 
     * <p> 
     * Note: This method is not synchronized. This way it is possible 
     * to get a value that is invalid already when this call returns. 
     * But there is no use in synchronizing it, because then 
     * you only move the point of interleaving _after_ this call. 
     * 
     * @return the maximum of active jobs that this pool can provide 
     */ 
    public int getMaxSize() { 

        return _maxSize; 
    } 

    /** 
     * Get the amount of jobs that are running now 
     * <p> 
     * Note: This method is not synchronized. This way it is possible 
     * to get a value that is invalid already when this call returns. 
     * But there is no use in synchronizing it, because then 
     * you only move the point of interleaving _after_ this call. 
     * 
     * @return the amount of jobs that are running now 
     */ 
    public int getCurrentActive() { 

        return _currentActive; 
    } 

    /** 
     * Schedule the given job for runtime. 
     * <p> 
     * The jobs are scheduled using the FIFO way, so first come, 
     * first scheduled. 
     * <p> 
     * This method returns instantly, instead of blocking when the 
     * maximum amount of running threads is reached. 
     * 
     * @param p_job the job to be scheduled 
     */ 
    synchronized 
    public void runJob( final Runnable p_job ) { 

        final PooledThread l_thread; 

        _waitingJobs.addLast( p_job ); 

        // We check if we can start processing jobs right now 
        // 
        if ( !(getMaxSize() < getCurrentActive()) ) { 

            // There are no threads waiting in the pool 
            // but we need to run a job, so we create a new one 
            // 
            if ( _inactiveThreads.isEmpty() ) { 

                l_thread = new PooledThread(); 
                l_thread.start(); 
            } else 
                l_thread = (PooledThread) _inactiveThreads.pop(); 

            // Run the first job found in the _waitingJobs queue 
            // 
            l_thread.wakeAndRunJob( (Runnable) _waitingJobs.removeFirst() ); 
            increaseActive(); 
        } 
    } 

    /** 
     * This method is used by threads that run jobs when the job is done. 
     * The thread is then either pushed on the _inactiveThreads Stack or 
     * told that he can gracefully die 
     * 
     * @param p_finishedThread the thread that is saying that he is done 
     * @return true if the Thread must keep on running, otherwise it must 
     * exit it's run method and die 
     */ 
    synchronized 
    private boolean push( final PooledThread p_finishedThread ) { 

        boolean l_letTreadLive = true; 

        // The thread already stopped being active 
        decreaseActive(); 

        // Check if the Thread may die. When we have enough threads waiting and 
        // being active, then we can let this die 
        if ( getMaxSize() <= _inactiveThreads.size() + getCurrentActive() ) { 
            l_letTreadLive = false; 
        } 

        // If there less jobs running then the maximum, then we must check if we 
        // can start schedule some 
        // 
        if ( getMaxSize() > getCurrentActive() ) { 

            // There are jobs waiting. Schedule one 
            // 
            if ( !_waitingJobs.isEmpty() ) { 
                p_finishedThread.wakeAndRunJob( (Runnable) _waitingJobs.removeFirst() ); 
                increaseActive(); 
            } else 
                _inactiveThreads.push( p_finishedThread ); 
        } 

        return l_letTreadLive; 
    } 

    /** 
     * Initialize the pool with initialSize Threads 
     * Not synchronized because it is intented only to be called 
     * by the constructor 
     * @param p_initialSize 
     */ 
    private void initPool( final int p_initialSize ) { 

        PooledThread l_thread; 

        for ( int i = 0; i < p_initialSize; ++i ) { 
            l_thread = new PooledThread(); 
            // Let the VM die when these Threads are the only one left 
            l_thread.setDaemon( true ); 

            l_thread.start(); 
            _inactiveThreads.push( l_thread ); 
        } 
    } 

    /** 
     * Increase the count for active jobs 
     * No need to synchronize it, because it is only called 
     * in synchronized methods 
     */ 
    private void increaseActive() { 

        ++_currentActive; 

        if ( _currentActive == getMaxSize() ) 
            firePoolFullEvent(); 
    } 

    /** 
     * Decrease the count for active jobs 
     * No need to synchronize it, because it is only called 
     * in synchronized methods 
     */ 
    private void decreaseActive() { 
        --_currentActive; 

        if ( _currentActive == 0 && _waitingJobs.isEmpty() ) 
            firePoolEmptyEvent(); 
    } 

    /** 
     * The internal Threads used by <code>FIFOThreadPool</code>. 
     * <p> 
     * The class <code>PooledThread</code> is a Thread that does not 
     * stop after the given runnable is finished with it's run method, 
     * but instead asks the Threadpool if there is another job. 
     * If not, this thread goes to sleep until a wakeAndRunJob( ) is called 
     */ 
    private class PooledThread extends Thread { 

        // Our current job 
        // 
        private Runnable _job = null; 

        /** 
         * Construct a PooledThread with no current job 
         */ 
        public PooledThread() { 

            this( null ); 
        } 

        /** 
         * Construct a PooledThread with a given job as current job 
         * @param p_job the current job 
         */ 
        public PooledThread( final Runnable p_job ) { 

            setJob( p_job ); 
        } 

        /** 
         * Set the current job to p_job 
         * @param p_job the new current job 
         */ 
        private void setJob( final Runnable p_job ) { 

            _job = p_job; 
        } 

        /** 
         * Set the given job as new current job and wake this Thread. 
         * <p> 
         * NOTE: This method is synchronized to make sure that it 
         * only can be called when this state is in the wait() state 
         * (caused by run) and that a second call on this method does not 
         * create skipped jobs 
         * 
         * @param p_job the new current job 
         */ 
        synchronized 
        public void wakeAndRunJob( final Runnable p_job ) { 

            setJob( p_job ); 
            notify(); 
        } 

        /** 
         * This run method is so modified that it doesn't stop until 
         * a <code>ThreadPool</code> says so. 
         * In it's run it checks if it has a job to run. If not, 
         * then this thread goes to sleep until waked by wakeAndRunJob() 
         * or Interrupted. If it has a job, it will run it and then 
         * disposes the job. 
         * <p> 
         * NOTE: This method is synchronized to make sure that it can't
         * be interleaved with wakeAndRunJob. That would cause really 
         * unpredictable results. (imagine a insertation of a job after 
         * run finished the running of the previous job) 
         */ 
        synchronized 
        public void run() { 

            boolean l_stop = false; 

            while ( !l_stop ) { 

                // If there is no job to do, go to sleep until notifyed 
                // by wakeAndRunJob 
                // 
                if ( _job == null ) 
                    try { 

                        wait(); 
                    } catch ( InterruptedException ie ) { 

                        ie.printStackTrace(); 
                        continue; 
                    } 

                // If the thread woke up and has data 
                // 
                if ( _job != null ) 
                    _job.run(); 

                setJob( null ); 

                // Check if we need to be put back on the stack 
                // or that we can die gracefully 
                // 
                l_stop = !push( this ); 
            } 
        } 
    } 
}


BlockedThreadedQuickSort
Ik heb ervoor gekozen om een blocking en een non-blocking versie te maken van de threadedquicksort. De blocking versie is waar het allemaal hier om ging (en die ik uiteindelijk niet gebruikt heb) post ik hier.
Hij maakt gebruik van een Semafoor om te blocken en class waar hij van extends is niet veel anders dan degene in mijn startpost
Java:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
public class BlockingThreadedQuickSort extends ThreadedQuickSort {

    // Doug Lea's Semaphore used for blocking. It is initialized by zero so
    // that a aquire() will result in a blocking Thread
    //
    private final Semaphore _finished = new Semaphore( 0 );

    /**
     * Construct a BlockingThreadedQuickSort with a given maximum of active Threads
     * @see ThreadedQuickSort#ThreadedQuickSort
     * @param p_maxSize maximum of active Threads
     */
    public BlockingThreadedQuickSort( final int p_maxSize ) {

        super( p_maxSize );
        getPool().setPoolListener( new JobsDoneListener() );
    }

    /**
     * Sort the given p_length elements of p_source from index p_offset and up by the quicksort algorithm.
     * <p>
     * Note that this sortingmethod is entirely equal to <code>QuickSort</code> except that it uses threads for
     * each recurrent call.
     *
     * @see QuickSort#quicksort
     *
     * @param p_source the input array that needs to be sorted
     * @param p_offset the offset boundary
     * @param p_length the amount of elements to sort
     */
    public void quicksort( final int[] p_source, final int p_offset, final int p_length ) {

        super.quicksort( p_source, p_offset, p_length );

        /*
         * Try to aquire() on the Semaphore. This will only be possible if
         * our JobsDoneListener has called to Semaphore.release. If the release
         * is still not done, then this will result in a blocking thread.
         */
        try {
            _finished.acquire();
        } catch ( InterruptedException l_interrupt ) {
            l_interrupt.printStackTrace();
        }
    }

    /**
     * A PoolListener that puts the Semaphore up (release()) when he gets the message that
     * it is empty (impling that the sorting is done, because this class uses it's own ThreadPool)
     */
    private class JobsDoneListener implements PoolListener {

        // If the pool is empty, then the jobs must be done.
        //
        public void poolEmpty( final PoolEvent p_event ) {

            _finished.release();
        }

        // We don't care if the pool has reached it's maximum amount of threads
        //
        public void poolFull( final PoolEvent p_event ) {
        }
    }
}


Iedereen hartelijk bedankt voor de suggesties en voorstellen :)
Als men nog kritiek heeft, spui alstublieft :)

[ Voor 3% gewijzigd door Glimi op 08-10-2003 22:48 ]


Verwijderd

glimi, ik moet zeggen dat ik je oplossing om threads die klaar zijn met een job via de push() methode gelijk een nieuwe job te laten ophalen erg elegant vind. Je hebt nu namelijk geen controller thread meer nodig in je threadpool...
Ik ga die oplossing gelijk ff implementeren in m'n eigen classes (das dan wel c++ maar je gekozen designoplossing werkt daar ook), thanks!

edit : glimi, performt deze methode nou ook beter dan een standaard quicksort?

[ Voor 11% gewijzigd door Verwijderd op 09-10-2003 15:19 ]


  • Glimi
  • Registratie: Augustus 2000
  • Niet online

Glimi

Designer Drugs

Topicstarter
(overleden)
Verwijderd schreef op 09 October 2003 @ 15:08:
glimi, ik moet zeggen dat ik je oplossing om threads die klaar zijn met een job via de push() methode gelijk een nieuwe job te laten ophalen erg elegant vind. Je hebt nu namelijk geen controller thread meer nodig in je threadpool...
Ik ga die oplossing gelijk ff implementeren in m'n eigen classes (das dan wel c++ maar je gekozen designoplossing werkt daar ook), thanks!
Kijk wel uit naar de 'kunstgreep' bij het setten van een nieuwe maxSize. Immers het kan voorkomen dat men 0 of minder opgeeft bij maxSize, waardoor alle gecachde threads doodgaan, terwijl er misschien nog jobs uitgevoerd moeten worden. Je moet er dus voor zorgen dat als er een nieuwe maxSize geset wordt, er weer aangevangen wordt met jobs processen. Dit is in mijn code een beetje gaar gedaan door de laatste job uit de queue te pakken en die te gebruiken bij een runJob() call.
Had een beetje tijdnood
edit : glimi, performt deze methode nou ook beter dan een standaard quicksort?
Geen idee, maar stel ik zou een eerlijke vergelijking maken (dwz creatie van het ThreadedQuickSort object en de Threadpool en dan nog sorten) met Arrays.sort(), dan vrees ik dat de array-grootte wel heel groot mag zijn om de initialisatie tegen te gaan.
Helaas heb ik geen manier om de boel te testen, ik werk hier met een 1 CPU systeem.
Als iemand wil testen, dan kan hij hier de sources vinden
Pagina: 1