【问题标题】:How to use multithreading in a loop in Java如何在 Java 中的循环中使用多线程
【发布时间】:2018-02-21 21:08:14
【问题描述】:

这就是我想要做的。我在while loop 中记录来自不同传感器的数据,直到用户停止记录。我想每秒记录尽可能多的数据。传感器需要不同的时间来返回一个值,介于 200 毫秒和 3 秒之间。因此,依次调用传感器不是一种选择。

依次调用传感器如下所示:

List<DataRow> dataRows= new ArrayList<DataRow>();

while (recording) {
   DataRow dataRow = new DataRow();

   dataRow.setDataA(sensorA.readData());
   dataRow.setDataB(sensorB.readData());
   dataRow.setDataC(sensorC.readData());

   dataRows.add(dataRow);
}

根据传感器,读取数据看起来(非常简化)

public class SensorA {

   public SensorAData readData(){
      sensorA.startSensing();

      try {
          TimeUnit.MILLISECONDS.sleep(750);
      } catch (InterruptedException e) {
          Thread.currentThread().interrupt(); 
      }

      return sensorA.readAndConvertByteStream();
   }
}

要利用多线程,SensorA 能否在循环中实现Callable 和接收Future 对象?还是应该将while loop 放在实现接口Runnablerun() 方法中?

基本上,即使循环已经至少进一步迭代,Java(或线程)是否可以写入正确的dataRow 对象?如果不是,如何解决这个问题?

【问题讨论】:

  • 您可以使用java.util.concurrent 包中的线程安全Collection 并将其传递给所有传感器线程。传感器线程从ThreadPoolExecutor 分配。主线程什么也不做,只是在对象上wait() 或重复Thread.sleep(),直到抛出InterruptedException,此时你关闭ThreadPoolExecutor
  • Re: Thread.currentThread().interrupt(); 如果你只是要退出我认为不设置中断位是可以接受的。您只需在代码的其他部分需要检测中断位并因此退出时进行设置。这意味着如果readAndConvertBytes 看到中断位,它可能会提前退出。可能不是你想要的。
  • can Java (or a thread) write to the correct dataRow object 这在很大程度上取决于您没有向我们展示的细节。读/写循环必须是线程安全的,我们无法判断它们是否是线程安全的。如果对象不与任何其他线程共享,那么它可能没问题(但要注意在底层共享全局状态的对象)。
  • IMO 您需要更加仔细地考虑您的要求。编写一个每秒轮询传感器 A 五次的循环和每三秒轮询传感器 B 一次的不同循环以及在不同线程中运行这些循环很容易。这与最初发明线程的目的非常接近。但是,我无法通过阅读您的问题来判断您想要 做什么 以非常不同的速度传入的这两个数据流。您希望每秒输出多少数据行?您想在哪一行查看哪些数据?

标签: java multithreading loops runnable callable


【解决方案1】:

如果我正确理解您的需求,这可能是您想要的解决方案:

  • 在每次迭代中,n 个传感器由 n 个并发线程读取,
  • 如果所有线程都收集了传感器数据,则将新结果行添加到列表中

工作代码:

public class TestX {

    private final ExecutorService pool = Executors.newFixedThreadPool(3);
    private final int N = 10;

    // all sensors are read sequentially and put in one row
    public void testSequential() {
        int total = 0;
        long t = System.currentTimeMillis();

        for (int i = 0; i < N; i++) {
            System.out.println("starting iteration " + i);

            int v1 = getSensorA();    // run in main thread
            int v2 = getSensorB();    // run in main thread
            int v3 = getSensorC();    // run in main thread

            // collection.add( record(v1, v2, v3)
            total += v1 + v2 + v3;
        }

        System.out.println("total = " + total + "   time = " + (System.currentTimeMillis() - t) + " ms");
    }

    // all sensors are read concurrently and then put in one row
    public void testParallel() throws ExecutionException, InterruptedException {
        int total = 0;
        long t = System.currentTimeMillis();

        final SensorCallable s1 = new SensorCallable(1);
        final SensorCallable s2 = new SensorCallable(3);
        final SensorCallable s3 = new SensorCallable(3);

        for (int i = 0; i < N; i++) {
            System.out.println("starting iteration " + i);

            Future<Integer> future1 = pool.submit(s1);  // run in thread 1
            Future<Integer> future2 = pool.submit(s2);  // run in thread 2
            Future<Integer> future3 = pool.submit(s3);  // run in thread 3

            int v1 = future1.get();
            int v2 = future2.get();
            int v3 = future3.get();

            // collection.add( record(v1, v2, v3)
            total += v1 + v2 + v3;
        }

        System.out.println("total = " + total + "   time = " + (System.currentTimeMillis() - t) + " ms");
    }

    private class SensorCallable implements Callable<Integer> {

        private final int sensorId;

        private SensorCallable(int sensorId) {
            this.sensorId = sensorId;
        }

        @Override
        public Integer call() throws Exception {
            switch (sensorId) {
                case 1: return getSensorA();
                case 2: return getSensorB();
                case 3: return getSensorC();
                default:
                    throw new IllegalArgumentException("Unknown sensor id: " + sensorId);
            }
        }
    }

    private int getSensorA() {
        sleep(700);
        return 1;
    }

    private int getSensorB() {
        sleep(500);
        return 2;
    }

    private int getSensorC() {
        sleep(900);
        return 2;
    }

    private void sleep(long ms) {
        try {
            Thread.sleep(ms);
        } catch (InterruptedException e) {
            // ignore
        }
    }

    public static void main(String[] args) throws ExecutionException, InterruptedException {
        new TestX().testSequential();
        new TestX().testParallel();
    }
}

和输出:

starting iteration 0
starting iteration 1
starting iteration 2
starting iteration 3
starting iteration 4
starting iteration 5
starting iteration 6
starting iteration 7
starting iteration 8
starting iteration 9
total = 50   time = 21014 ms

starting iteration 0
starting iteration 1
starting iteration 2
starting iteration 3
starting iteration 4
starting iteration 5
starting iteration 6
starting iteration 7
starting iteration 8
starting iteration 9
total = 50   time = 9009 ms

-- 编辑--

在 java 8 中,您可以使用方法引用来摆脱 Callable 类,只需编写:

Future<Integer> future1 = pool.submit( this::getSensorA() );
Future<Integer> future2 = pool.submit( this::getSensorB() );
Future<Integer> future3 = pool.submit( this::getSensorC() );

【讨论】:

  • 感谢@przemek-hertel 的精心提议。据我了解future.get() 将等待结果并有效地暂停循环,直到最慢的传感器准备好。有没有办法继续循环,一准备好就写入慢速传感器的数据?
  • 是的,有可能。您可以将生产者-消费者模式与 BlockingQueue 一起使用。您的循环是生产者(主线程)并将请求对象放入队列。您必须启动从队列中获取请求的消费者线程池,读取传感器(在请求中描述),最后将结果放入结果集合。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-03-09
  • 2017-09-25
  • 1970-01-01
相关资源
最近更新 更多