Wie mehreren Threads Aufgaben zuweisen und abarbeiten lassen?

P@

Mitglied
Hallo allerseits,

bin neu hier im Forum. Hab schon einige Erfahrungen mit Java gemacht, beiße mir aber am Thema Threads die Zähne aus. Während ich mich u.a. durch die Java-Insel und einige Tutorials gequält habe, sind schon etliche Programmierversuche ins Land gegangen. Meine Versuche mit synchronized, volatile, wait, notify und Co. haben bei mir nur dazu geführt, das mehrere Threads nacheinander und nicht parallel von einem Hauptthread bedient wurden. Da habe ich wohl was nicht richtig verstanden.
In dem folgenden Programm habe ich einfach mal drauf losgelegt und das Ergebnis scheint so zu funktionieren, wie ich's mir vorstelle: parallel und - ich hoffe - fehlerfrei.

Im s.g. Hauptthread (Ht) werden zwei Threads (Consumer) erzeugt und mit Aufgaben (Zahlen) gefüttert. Hat ein Thread eine Aufgabe erledigt, durchläuft er ständig eine Schleife, um zu prüfen, ob der Ht eine weitere Aufgabe für ihn hat. Nach erfolgreicher Bearbeitung wird, als Signal erneuter Aufnahmebereitschaft an den Ht, ein Trigger auf false gesetzt. Dieser wird zyklisch vom Ht abgefragt. Erkennt dieser den Trigger als false, schickt er erneut eine Aufgabe an den jeweiligen Thread und setzt dessen Trigger auf true, was für den Consumer wiederum das Signal zum erneuten run-Durchlauf ist.
Wenn alle Aufgaben erledigt sind, erhält jeder Consumer-Thread ein Abbruchsignal (-1 als Aufgabe), um sich selber zu beenden.

Java:
public class Consumer extends Thread{
	private String name;
	private boolean trigger;
	private int aufgabe; //Was der Thread machen soll
	
	public Consumer( String n ){
		this.name= n;
		this.trigger= false;
		this.aufgabe= 0;
	}
	
	public void setAufgabe( int a ){
		this.aufgabe= a;
	}
	
	public boolean getTrigger(){
		return this.trigger;
	}
	
	public void setTrigger(){
		this.trigger= true;
	}
	
	public void run(){
		while( true ){
			while( !trigger ){
				try { sleep( 1 ); } catch ( InterruptedException e ){ e.printStackTrace(); }
			}
			/********** Aufgabenbereich *****************************************************/
			
			if( aufgabe== -1 ){ // -1 als Abbruchsignal, Thread beenden
				System.out.println("\n" + name + " beendet" );
				break;
			}
			System.out.println( name + "  " + aufgabe );
			try { sleep( aufgabe ); } catch ( InterruptedException e ){ e.printStackTrace(); }
			
			/********************************************************************************/
			
			trigger= false;
		}
	}
	
	public static void main(String[] args) {
		int[] aufgaben= { 10, 20, 30, 40, 500, 60, 70, 80, 90, 100, 110, 120, 130, 140, 150, 160, 170, 180, 190, 200 };
		int threadanz= 2;
		
		/***** erzeugen und starten von 2 (threadanz) Consumer-Threads *************/
		Consumer []c= new Consumer[threadanz];
		for( int i= 0; i< threadanz; i++ ){
			c[i]= new Consumer("c" + i);
			c[i].start();
		}
		
		/***** abarbeiten des aufgaben-Arrays ***************/
		int i= 0;
		int j= 0;
		while( i< aufgaben.length ){
			if( j>= threadanz ) j= 0;
			if( !c[j].getTrigger() ){
				c[j].setAufgabe(aufgaben[i]);
				c[j].setTrigger();
				i++;
			}
			j++;
		}
		
		/***** alle Consumer-Threads erhalten das Signal zum Beenden (-1) ****/
		i= 0;
		while( i< c.length ){
			if( i>= c.length ) i= 0;
			if( !c[i].getTrigger() ){
				c[i].setAufgabe(-1);
				//c[i].interrupt();
				c[i].setTrigger();
				i++;
			}
		}
	}
}

Meine eigentliche Frage bzw. Bitte lautet nun, mir die richtige Richtung zu weisen. Sprich: Hat jemand einen besseren Lösungsvorschlag, um mehrere Threads bestimmte Aufgaben parallel abarbeiten zu lassen. Ich bitte Euch um aktive Meinungsäußerung.
Und noch etwas. Wie ich bereits erwähnte, bin ich neu hier in diesem Forum. Das soll heißen, dass man mir, falls dieses Threads-füttern-Thema hier bereits zur Genüge gefragt, diskutiert und behandelt wurde, vergeben möge. Gleiches gilt für den Fall, dass ich zu :autsch: sein sollte, nach dem richtigen Foreneintrag zu suchen.
Ich bin für jede Hilfe und jeden Hinweis dankbar.

MfG Patrick
 
Zuletzt bearbeitet:
Hm, ich verstehe dein Problem nicht. Dein Programm erzeugt bei mir den output unten und da kann man erkennen, dass sich die threads schön brav abwechseln. Allerdings starte ich das aus Eclipse und wenn du es direkt auf eine VM startest könnte das schon ein verändertes Taskswitching bedeuten.

Ich vermute der Code macht das was du wolltest.

Code:
c1  20
c0  10
c0  30
c1  40
c0  500
c1  60
c1  70
c1  80
c1  90
c1  100
c1  110
c0  120
c1  130
c0  140
c1  150
c0  160
c1  170
c0  180
c1  190
c0  200

c1 beendet

c0 beendet
 
Hm, ich verstehe dein Problem nicht. Dein Programm erzeugt bei mir den output unten und da kann man erkennen, dass sich die threads schön brav abwechseln. Allerdings starte ich das aus Eclipse und wenn du es direkt auf eine VM startest könnte das schon ein verändertes Taskswitching bedeuten.

Ich vermute der Code macht das was du wolltest.

Code:
c1  20
c0  10
c0  30
c1  40
c0  500
c1  60
c1  70
c1  80
c1  90
c1  100
c1  110
c0  120
c1  130
c0  140
c1  150
c0  160
c1  170
c0  180
c1  190
c0  200

c1 beendet

c0 beendet

Hallo Andi,

exakt, gerade ab
Code:
c0 500
lässt sich sehr schön erkennen, dass c1 weiterhin "gefüttert" wird, während c0 noch beschäftigt ist.
Ich programmiere auch mit Eclipse. Aber wegen Deiner Bedenken bezüglich des veränderten Taskswitchings kann ich sagen, dass das Programm als jar-Datei, unter Windows und Linux die gleiche Ausgabe wie unter Eclipse erzeugt, falls Du das gemeint haben solltest.

MfG Patrick
 
Producer-Consumer Probleme löst man z.B. mit Mutex aus dem concurrency package. Es gibt da noch mehr tolle Implementierungen, die das Handling mit Threads nicht mehr nötig machen.

Concurrent Programming with J2SE 5.0

Hallo FArt,
vielen Dank für Deine Antwort inkl. dem Link. Ja, diese "lustigen" Semaphoren. Was das ist, und was es macht, habe ich schon begriffen. Leider hapert's da bei mir an der Umsetzung. Wie ich bereits erwähnte, haben meine bisherigen Versuche mit synchronized & Co. - und dazu gehört auch leider der Semaphoren-Ansatz - nur dazu geführt, dass der Hauptthread immer warten musste, bis die Consumer ihre Locks wieder freigegeben haben, wodurch sie lediglich seriell und nicht parallel abgearbeitet wurden. Ich hoffe, ich drücke mich klar genug aus. Deshalb galt ja meine eigentliche Bitte dem Einbringen von Vorschlägen, Codeschnipseln, wie man dies z.B. mittels Semaphoren parallelisieren kann.

MfG Patrick

PS: Kannst Du mir bitte erklären, wie du auf den Link zu dem Artikel bei sun.com gekommen bist?
 
Zuletzt bearbeitet:
Nachtrag:

Ein entscheidendes Problem ist, dass wenn die Consumer-Threads ziemlich lange mit einer Aufgabe beschäftigt sind, der Hauptthread sich in einer Art Endlosschleife beim Abfragen der ist-bereit-Zustände der Consumerthreads (Trigger) tot macht; deutlich sichtbar an der Prozessorauslastung. Dafür bräuchte ich eine optimierte Lösung.

Ich möchte an dieser Stelle nochmals allen danken, die sich bisher für mein Problem interessiert haben.

MfG Patrick
 
Zuletzt bearbeitet:
Schau dir mal an wie ich das ganze (ungefähr) umsetzen würde:

Java:
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;


class Consumer implements Runnable {
	
	private final String name;
	private final BlockingQueue<Integer> queue;
	
	public Consumer(String name, BlockingQueue<Integer> queue) {
		this.name = name;
		this.queue = queue;
	}
	
	@Override
	public void run() {
		try {
			while(true) {
				// Aufgabe aus Queue entnehmen, take blockiert falls Queue leer ist
				Integer task = queue.take();
				
				// Aufgabe -1 als Abbruchsignal
				if(task == -1) {
					break;
				}
				
				// Eigentliche Aufgabe ausführen
				System.out.println(name + " Aufgabe " + task);
				Thread.sleep(task);
			}
		} catch(InterruptedException e) {
			e.printStackTrace();
		}
		System.out.println(name + " beendet");
	}
}

public class Test {
	
	public static void main(String[] args) {
        int[] tasks = { 10, 20, 30, 40, 500, 60, 70, 80, 90, 100, 110, 120, 130, 140, 150, 160, 170, 180, 190, 200 };
        // Maximal 10 Aufgaben in Queue halten
        BlockingQueue<Integer> queue = new LinkedBlockingQueue<Integer>(10);        
        
        // Maximal 2 Threads verwenden
        int numberOfThreads = 2;
        ExecutorService executor = Executors.newFixedThreadPool(numberOfThreads);
        
        // 2 Consumer ausführen
        for(int i = 0; i < numberOfThreads; i++) {
        	executor.execute(new Consumer("c" + i, queue));
        }
        
        try {
        	// Aufgaben nacheinander in Queue einfügen, put blockiert falls Queue voll ist
	        for(int i = 0; i < tasks.length; i++) {
	        	queue.put(tasks[i]);
	        }
	        // Abbruchsignal in Queue einfügen
	        for(int i = 0; i < numberOfThreads; i++) {
	        	queue.put(-1);
	        }
        } catch(InterruptedException e) {
        	e.printStackTrace();
        }
        
        // Alle Threads beenden
        executor.shutdown();
	}
}
 
Nachtrag:

Ein entscheidendes Problem ist, dass wenn die Consumer-Threads ziemlich lange mit einer Aufgabe beschäftigt sind, der Hauptthread sich in einer Art Endlosschleife beim Abfragen der ist-bereit-Zustände der Consumerthreads (Trigger) tot macht; deutlich sichtbar an der Prozessorauslastung. Dafür bräuchte ich eine optimierte Lösung.
...

du lässt deine HT aber auch nicht viel Zeit^^
Java:
while (!trigger) {
				try {
					sleep(1);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
			}
sleep(1) ist 1ms, deine eigentlichen Tasks legst du aber mindestens 10, im Extremfall 500 ms schlafen...evtl. kannst du da dynamisch eine Variable setzen, die in Abhängigkeit der Aufgabenlänge steht. Ansonsten würde ich an deiner stelle diesen Wert einfach erhöhen.
 
Ich glaub der andiv hat's. Ich werd verrückt. Alles läuft parallel, und der Hauptthread macht sich nicht mehr tot, wenn die Consumer etwas länger mit ihren Aufgaben beschäftigt sind. Menschenskind, beim Barte Odins, Heureka und fettes, fettes Merci Andi (entschuldigung, aber ich freu mich eben). Werd mir den Code glaub ich einrahmen lassen.
Kurze Anmerkung: Es funktioniert auch, wenn man queue mit einer Größe von 1 initialisiert &
darf ich fragen, ob Du mir weiterführende Literatur empfehlen kannst?

MfG Patrick
 
Zuletzt bearbeitet:
du lässt deine HT aber auch nicht viel Zeit^^
Java:
while (!trigger) {
				try {
					sleep(1);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
			}
sleep(1) ist 1ms, deine eigentlichen Tasks legst du aber mindestens 10, im Extremfall 500 ms schlafen...evtl. kannst du da dynamisch eine Variable setzen, die in Abhängigkeit der Aufgabenlänge steht. Ansonsten würde ich an deiner stelle diesen Wert einfach erhöhen.

Diese Idee ist mir auch schon gekommen. Allerdings sollen die Consumer im fertigen Progi nicht einfach nur eine best. Zeit lang schlafen, sondern Bilder von png in jpg umwandeln. Man könnte zwar z.B. in Abhängigkeit der Dateigröße eine dynamische Variable schätzen, aber mehr eben nicht.

Aber Danke für den Hinweis
 
Bei der Semaphorenimplementierung hier wartet nur der Consumer, und zwar bis es etwas zu tun gibt. Der Producer kommt sofort zurück.
Den Link findet man, wenn man sinnvolle Suchbegriffe hat (Semaphore z.B.) und das mit Begriffen wie java und tutorial verbindet.
Sonst: mal stöbern, was Oracle so alles an Tutorials, Dokus und technischen Artikeln so bietet.
 
Genaugenommen brauchst du die Größenbeschränkung der Queue hier gar nicht, aber sie kann nützlich sein wenn man will dass der Aufgabenberg eine bestimmte Größe nicht überschreitet. Auch den ExecutorService brauchst du hier nicht unbedingt, man hätte genauso gut von Hand zwei Threads erstellen können. Aber so hast du wenigstens mal gesehen wie man die verschiedenen Klassen bei einem Producer-Consumer-Problem anwenden kann.

Was Literatur angeht kann ich allgemein zu fortgeschrittenen Java-Themen "Effective Java" und zu Multithreading "Java Concurrency in Practice" empfehlen.
 

Zurück
Oben