Code 'innerhalb' des synchronen Bereichs einer BlockingQueue ausfuehren..?

sirbender

Top Contributor
Hallo,
eine BlockingQueue ist ja threadsafe bzgl. Operationen wie add() und take().
Nun wuerde ich fuer einige Experimente die ich damit mache gerne ausgeben koennen wie gross z.B. die Queue ist usw. nachdem ich z.B. ein add() ausgefuehrt habe.

Leider ist das alles nicht 'atomic'. Soll heissen, zwischen meinem add() und System.out.println(queue.size()) kann eine Menge passieren (eine anderer Thread fuegt der Queue Elemente hinzu bzw. entfernt diese). Ich weiss also nicht wie gross die Queue unmittelbar vor bzw. nach dem add() war.

Kann ich irgendwie Code 'uebergeben' der vor bzw. nach dem add() ausgefuehrt wird, allerdings innerhalb des 'Locks' des add() Aufrufs?

Bzw. gibt es eine andere Loesung fuer dieses Problem?
 
Ich würde mal behaupten, das wird so nicht möglich sein.

uU Decorator-Pattern mit eigener Synchronisation nutzen?
 
Ich hab mal kurz ein bischen minimalen Testcode geschrieben auf dessen Grundlage man diskutieren kann.
Wie hast du dir das mit dem Decorator-Pattern und eigener Synchronisation vorgestellt?

Ich habe im Code markiert welchen Bereich ich synchron ausfuehren will.

Vielen Dank!


Java:
package exp;

import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;

/**
 * Usage example, based on a typical producer-consumer scenario. Note that a
 * BlockingQueue can safely be used with multiple producers and multiple
 * consumers.
 */
public class BlockingCueExample {
 
    static abstract class Producer<T> implements Runnable {
      private final BlockingQueue<T> queue;
      Producer(BlockingQueue<T> q) { queue = q; }
      public void run() {
        try {
                while (true) {
                    T produce = produce();
                   
                    // Wie kann ich sicher gehen, dass dieser Bereich als Block ausgefuehrt wird und zwischendrin kein Thread auf die Queue zugreift?

                    // Anfang
                   
                    System.out.println("p: " + produce + "[" + Thread.currentThread() + "]");
                    System.out.println(queue.size());
                    queue.put(produce);
                    System.out.println("p: " + produce + "[" + Thread.currentThread() + "]");
                    System.out.println(queue.size());
                   
                    // Ende
                }
        } catch (InterruptedException ex) { System.out.println("... handle ...");}
      }
      abstract T produce() throws InterruptedException;
    }

    static class Consumer<T> implements Runnable {
      private final BlockingQueue<T> queue;
      Consumer(BlockingQueue<T> q) { queue = q; }
      public void run() {
        try {
          while (true) { consume(queue.take()); }
        } catch (InterruptedException ex) { System.out.println("... handle ..."); }
      }
      void consume(T x) { System.out.println("c: " + x + "[" + Thread.currentThread() + "]"); }
    }

     
    public static void main(String[] args) {
        BlockingQueue<String> q = new LinkedBlockingQueue<String>();
        Producer<String> p1 = getProducer(q);
        Producer<String> p2 = getProducer(q);
        Consumer<String> c1 = new Consumer<String>(q);
        Consumer<String> c2 = new Consumer<String>(q);
        new Thread(p1).start();
        new Thread(p2).start();
        new Thread(c1).start();
        new Thread(c2).start();   
    }

    private static Producer<String> getProducer(BlockingQueue<String> q) {
        return new Producer<String>(q) {
            @Override
            String produce() throws InterruptedException {
                long timeMillis = System.nanoTime();
                Thread.sleep(500);
                return "" + timeMillis;
            }
        };
    }
}
 
Java:
      public void run() {
        try {
                while (true) {
                    T produce = produce();
                
                    // Wie kann ich sicher gehen, dass dieser Bereich als Block ausgefuehrt wird und zwischendrin kein Thread auf die Queue zugreift?

                   
                   synchronized(queue) {
                
                             System.out.println("p: " + produce + "[" + Thread.currentThread() + "]");
                             System.out.println(queue.size());
                             queue.put(produce);
                             System.out.println("p: " + produce + "[" + Thread.currentThread() + "]");
                             System.out.println(queue.size());
                       }
                
                    // Ende
                }
        } catch (InterruptedException ex) { System.out.println("... handle ...");}
      }
 
Und an der nächsten stelle vergisst man den synchronized-block und alles geht kaputt ^^


Ja, aber konkret mit dem aktuellen Beispiel von mir kann ich mir gerade nicht vorstellen wie die Sync. ablaufen muss?

Du implementierst selber BlockingQueue, in den Methoden synchronisierst du erst über ein eigenes Objekt und in dem Block machst du deine Ausgaben und delegierst weiter. Alternativ fügst du zusätzliche Methoden hinzu, denen du dann zusätzlich Consumer übergeben kannst
 
Ich hatte gehofft, dass es vielleicht eine Implementierung gibt, die fuer alle Methoden von BlockingQueue eine Methode mitdemselben Namen hinzufuegt und diese um zwei Runnables erweitert. Am Ende hat man dann also 2 Methoden und die eine ruft die andere auf mit Runnables auf Null gesetzt:
queue.add(T tobeAdded)
queue.add(T tobeAdded, Runnable executeBefore, Runnable executeAfter)

Ich hab es mir mal angeschaut und ich befuerchte, dass eine sichere Implementierung einer solchen BlockingQueue meine Faehigkeiten uebersteigt.

Waere aber denke ich eine gute Idee sowas zu haben.
 

Zurück
Oben