【问题标题】:Using a PriorityBlockingQueue to feed in logged objects for processing使用 PriorityBlockingQueue 输入记录的对象以进行处理
【发布时间】:2015-10-15 23:17:28
【问题描述】:

我有一个应用程序,它从多个序列化对象日志中读取对象并将它们交给另一个类进行处理。我的问题集中在如何高效、干净地读取对象并将它们发送出去。

代码是从旧版本的应用程序中提取的,但我们最终保持原样。直到上周才真正使用它,但我最近开始更仔细地查看代码以尝试改进它。

它打开 N 个ObjectInputStreams,并从每个流中读取一个对象并将它们存储在一个数组中(假设下面的inputStreams 只是对应于每个日志文件的ObjectInputStream 对象数组):

for (int i = 0; i < logObjects.length; i++) {
    if (inputStreams[i] == null) {
        continue;
    }
    try {
        if (logObjects[i] == null) {
            logObjects[i] = (LogObject) inputStreams[i].readObject();
        }
    } catch (final InvalidClassException e) {
        LOGGER.warn("Invalid object read from " + logFileList.get(i).getAbsolutePath(), e);
    } catch (final EOFException e) {
        inputStreams[i] = null;
    }
}

序列化到文件的对象是LogObject 对象。这是LogObject 类:

public class LogObject implements Serializable {

    private static final long serialVersionUID = -5686286252863178498L;

    private Object logObject;
    private long logTime;

    public LogObject(Object logObject) {
        this.logObject = logObject;
        this.logTime = System.currentTimeMillis();
    }

    public Object getLogObject() {
        return logObject;
    }

    public long getLogTime() {
        return logTime;
    }
}

一旦对象在数组中,它就会比较日志时间并发送最早时间的对象:

// handle the LogObject with the earliest log time
minTime = Long.MAX_VALUE;
for (int i = 0; i < logObjects.length; i++) {
    logObject = logObjects[i];
    if (logObject == null) {
        continue;
    }
    if (logObject.getLogTime() < minTime) {
        index = i;
        minTime = logObject.getLogTime();
    }
}

handler.handleOutput(logObjects[index].getLogObject());

我的第一个想法是为每个文件创建一个线程,该线程读取并将对象放入PriorityBlockingQueue(使用自定义比较器,该比较器使用LogObject 日志时间进行比较)。然后另一个线程可能会取出这些值并将它们发送出去。

这里唯一的问题是一个线程可以将一个对象放入队列并在另一个线程可以将一个可能有更早时间的对象放入队列之前将其取出。这就是为什么在检查日志时间之前首先读取对象并将其存储在数组中的原因。

此约束是否禁止我实施多线程设计?或者有什么方法可以调整我的解决方案以提高效率?

【问题讨论】:

    标签: java multithreading priority-queue


    【解决方案1】:

    据我了解,您需要严格按顺序处理LogObjects。在这种情况下,您的代码的初始部分是完全正确的。这段代码所做的是几个输入流的merge sort。您需要为每个流读取一个对象(这就是需要临时数组的原因),然后采用适当的(最小/最大)LogObject 并处理到处理器。

    根据您的上下文,您可能能够在多个线程中进行处理。您唯一需要更改的是将LogObjects 放入ArrayBlockingQueue 中,处理器可能会在多个独立线程上运行。另一种选择是发送LogObjects 以在ThreadPoolExecutor 中进行处理。最后一个选项更简单直接。

    但要注意路上的几个陷阱:

    • 要使该算法正常工作,必须对各个流进行排序。否则你的程序就坏了;
    • 当你在并行处理消息时,处理顺序严格来说是没有定义的。所以提出的算法只保证消息处理的开始顺序(调度顺序)。这可能不是您想要的。

    所以现在你应该面临几个问题:

    1. 真的需要处理订单吗?
    2. 如果需要,是否需要全局顺序(针对所有消息)或本地顺序(针对独立的消息组)?

    回答这些问题将对您进行并行处理的能力产生重大影响。

    如果第一个问题的答案是,遗憾的是,并行处理不是一个选项。

    【讨论】:

    • 我想问一下是否需要全局顺序。如果不是,那么我有一个简单的解决方案。如果是这样,我能改进吗?您提到了 ArrayBlockingQueue,但我不确定这是否是基于需要全局订单的改进。
    • 如果需要全局处理订单,那就无能为力了。除非您可以将处理分成几个步骤,否则某些步骤不需要订购。因此,只有那些步骤可以并行完成。很难说处理将如何加快。但是,如果您可以估计可以并行完成的工作的比例,那么可以使用 [Amdahl 定律][en.wikipedia.org/wiki/Amdahl%27s_law] 来确定可能的加速。
    【解决方案2】:

    我同意你的看法。扔掉它并使用PriorityBlockingQueue.

    这里唯一的问题是,如果线程 1 已经从文件 1 中读取了一个对象并将其放入队列中(并且文件 2 将要读取的对象具有更早的日志时间),则读取线程可以接受它并将其发送出去,从而产生一个稍后发送的日志对象

    这与平衡合并 (Knuth ACP vol 3) 的合并阶段完全相同。您必须从获得前一个最低元素的同一文件中读取下一个输入。

    此约束是否禁止我实现多线程设计?

    这不是一个约束。这是虚构的。

    或者有什么方法可以调整我的解决方案以提高效率?

    优先队列已经非常高效了。无论如何,您当然应该首先担心正确性。然后添加缓冲 ;-) 将 ObjectInputStreams 包裹在 BufferedInputStreams 周围,并确保在您的输出堆栈中有一个 BufferedOutputStream

    【讨论】:

    • 我知道它会占用日志时间最短的对象,但考虑一下:thread1 读入日志时间为 100 的object1 并将其放入队列中,@987654327 @ 将object1 从队列中取出并发送出去,thread2 读入日志时间为 50 的 object2 并将其放入队列中,等等。日志时间为 100 的对象在一个对象之前已发送50.
    • 这就是为什么我说最初所有对象都是从 N 个文件中读入的,然后在发送之前检查它们的日志时间。
    • 您能否解决我在第一条评论中概述的情况?一个线程可以将一个对象放入队列并在另一个线程放入可能有更早时间的对象之前将其取出。
    猜你喜欢
    • 2021-02-05
    • 1970-01-01
    • 2020-03-01
    • 1970-01-01
    • 1970-01-01
    • 2018-02-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多