【问题标题】:Concurrent Observer Pattern并发观察者模式
【发布时间】:2014-02-06 02:26:07
【问题描述】:

如果我遇到如下情况:

ObserverA、ObserverB、ObserverC 都继承自 AbstractObserver。

我创建了一个观察者列表:

List<AbstractObserver> list = new ArrayList<AbstractObserver>();
list.add(new ObserverA());
list.add(new ObserverB());
list.add(new ObserverC());

并且具有以下方法的某种处理程序在“MAIN”线程中运行:

public void eat(Food item) {
     for(AbstractObserver o : list) {
          o.eatFood(item);
     }
}

public void drink(Coffee cup) {
     for(AbstractObserver o : list) {
          o.drinkCoffee(cup);
     }
}

如何设计一个系统,让我可以在不同线程中运行观察者的每个eatFood 和drinkCoffee 方法?具体来说,当“MAIN”线程接收到事件(调用drink 或eat 方法)时,我将如何在它们自己的线程中运行ObserverA、ObserverB 和ObserverC 中的eatFood 或drinkCoffee 方法?

我希望为每个 AbstractObserver 子类实例设置不同的线程,因为目前我正在按顺序通知每个观察者,这可能会导致延迟。

【问题讨论】:

  • 我不明白作为观察者的对象如何在这里发挥作用。他们观察的主题是什么?它们与eatFood and DrinkCoffe 方法的作用有什么关系?
  • 不,不适合我。从我的角度来看,我不知道何时会调用吃喝(主线程方法)。他们可以相隔五分钟或更长时间被调用。我只希望每个单独的 AbstractObserver 子类实例能够在调用主线程的吃或喝方法时同时运行drinkCoffee 或eatFood。
  • 你对并发有什么限制吗?例如,可以在之前的 o.eatFood() 结束之前调用同一对象上的后续 o.eatFood() 吗?

标签: java multithreading concurrency observer-pattern


【解决方案1】:

我不是这方面的专家,但也许您可以使用生产者-消费者设置。在这里,作为被观察实体的生产者可以在其自己的线程中的队列上添加通知,消费者(此处的观察者)将从同一个队列中获得通知,但在其自己的线程中。

【讨论】:

    【解决方案2】:

    要详细说明 Hovercraft 的答案,观察者的基本实现可能如下所示:

    class ObserverA implements Runnable {
        private final BlockingQueue<Food> queue = new ArrayBlockingQueue<> ();
    
        public void eatFood(Food f) {
            queue.add(f);
        }
    
        public void run() {
            try {
                while (true) {
                    Food f = queue.take(); //blocks until some food is on the queue
                    //do what you have to do with that food
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                //and exit
            }
        }
    }
    

    因此,您调用 eatFood 的代码将立即从该方法返回,而不会阻塞您的主线程。

    您显然需要直接为观察者分配一个线程:new Thread(observerA).start(); 或通过 ExecutorService,这可能更容易且更可取。


    或者,您可以在“观察到的”对象级别创建线程:

    private static final ExecutorService fireObservers = Executors.newFixedThreadPool(10);
    
    public void eat(Food food) {
        for (AbstractObserver o : observers) {
            //(i)  if an observer gets stuck, the others can still make progress
            //(ii) if an observer throws an exception, a new thread will be created
            Future<?> f = fireObservers.submit(() -> o.dataChanged(food));
            fireObservers.submit(new Callable<Void>() {
                @Override public Void call() throws Exception {
                    try {
                        f.get(1, TimeUnit.SECONDS);
                    } catch (TimeoutException e) {
                        logger.warn("Slow observer {} has not processed food {} in one second", o, food);
                    } catch (ExecutionException e) {
                        logger.error("Observer " + o + " has thrown exception on food " + food, e.getCause());
                    }
                    return null;
                }
            });
        }
    }
    

    (我主要是从here 复制粘贴的——你可能需要根据自己的需要调整它)。

    【讨论】:

    • 谢谢,我想我明白这是如何工作的。不过我的问题是,如果我走这条路,我将不得不自己为 MAIN 线程中的每个方法创建队列。在您的示例中,我相信我还必须为 drinkFood 创建一个方法。在我的项目中,我有大约 20 种方法需要它,所以如果我尝试自己处理 20 个队列,然后在 run 方法中处理 queue.take() 方法,我担心它可能会变得笨拙。
    【解决方案3】:

    在这里做一些简化假设,即您不关心在吃/喝结束时收到通知,您还可以使用执行器框架将工作放入队列中:

      // declare the work queue
       private final Executor workQueue = Executors.newCachedThreadPool(); 
    
    
    
    
       // when you want to eat, schedule a bunch of 'eating' jobs
           public void eat(final Food item){
              for (final AbstractObserver o: list) {
                 workQueue.execute(new Runnable() {
    
                    @Override
                    public void run() {
                       o.eatFood(item); // runs in background thread
                    }
                 });
              }
           }
    

    退出程序时,必须关闭执行器:

       workQueue.shutdown();
    

    【讨论】:

    • 这个方法仍然需要你修改每个方法(吃,喝等),但至少你只有一个工作队列(不是20个),线程由执行器框架为你管理。如果您需要限制并发性,只需将执行程序换成使用 FixedThreadPool 的执行程序即可。
    • 我认为 assylias 的回答也是正确的,但这似乎是最简单的解决方案。
    猜你喜欢
    • 2016-02-20
    • 2023-04-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-11-22
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多