【问题标题】:Lazy sorting of entities in Java 8 Stream API on a daily basis?每天对 Java 8 Stream API 中的实体进行延迟排序?
【发布时间】:2017-09-20 03:23:33
【问题描述】:

我有一个大型 Java 8 流 (Stream<MyObject>),其中的对象如下所示:

class MyObject {
   private String string;
   private Date timestamp;

   // Getters and setter removed from brevity 
}

我知道第 1 天的所有时间戳都会在第 2 天之前到达,但在每一天中,时间戳可能会出现故障。我想每天使用 Stream API 对timestamp 订单中的MyObject 进行排序。由于 Stream 很大,我必须尽可能懒惰地执行此操作,即可以在内存中保存一天的 MyObject,但 不 可以保存更多比那个。

我怎样才能做到这一点?

2017-04-29 更新:

一个要求是我想在排序后继续处理同一个流!我想要这样的东西(伪代码):

Stream<MyObject> sortedStream = myStreamUnsorted().sort(onADailyBasis());

【问题讨论】:

  • 更多的是关于调度(何时做)或处理(如何做)的问题?你在使用 Spring 堆栈吗?
  • 我认为您很可能必须使用 Stream 的迭代器,以便在排序之前每天对元素进行分组。我不认为 Stream API 可以帮助你解决这种需求。
  • @AndrewTobilko 这是关于处理的,我没有使用 Spring。
  • 你的时间戳之间最小的 TimeUnit 是多少?我们是在谈论几秒钟还是更少?
  • 恐怕java流不太适合排序——尤其是元素很多的时候。

标签: java java-8 java-stream


【解决方案1】:

我建议以下解决方案:

将流的每个值存储在 TreeMap 中,以便立即对其进行排序。使用对象的时间戳作为键。

 Map<Date, MyObject> objectsOfTheDaySorted = new TreeMap<>();

我们需要知道最后必须从地图中删除哪个对象。它只是一个对象,但存储它的成员必须(有效)是最终的。所以我选择了一个简单的列表。

 List<MyObject> lastObject = new ArrayList<>();

将当前日期设置为整数。

 // just an example
 int currentDay = 23;

使用一个谓词来确定 currentDay 和任何经过对象的日期是否不匹配。

 Predicate<MyObject> predicate = myObject -> myObject.getTimestamp()
                    .toInstant()
                    .atZone(ZoneId.systemDefault())
                    .toLocalDate()
                    .getDayOfMonth() != currentDay;

现在直播您的直播。使用 peek() 两次。首先将对象放入地图中。 其次覆盖列表中的对象。 使用anyMatch()作为终端操作,之前的交 创建谓词。一旦出现第一个匹配的对象 从第二天开始,anyMatch() 终止流并返回 true。

 stream.peek(myObject -> objectsOfTheDaySorted.put(myObject.getTimestamp(), myObject))
       .peek(myObject -> lastObject.set(0, myObject))
       .anyMatch(predicate);

现在您只需要删除已经属于第二天的最后经过的对象,因此不属于您的地图。

 objectsOfTheDaySorted.remove(lastObject.get(0).getTimestamp());

完成。你有一个排序的对象地图,它们都只属于一天。希望这符合您的期望。请在下面找到一个块中的整个代码,以便一次更好地复制它。

 Map<Date, MyObject> objectsOfTheDaySorted = new TreeMap<>();
 List<MyObject> lastObject = new ArrayList<>();

 // just an example
 int currentDay = 23;

 Predicate<MyObject> predicate = myObject -> myObject.getTimestamp()
                    .toInstant()
                    .atZone(ZoneId.systemDefault())
                    .toLocalDate()
                    .getDayOfMonth() != currentDay;

 stream.peek(myObject -> objectsOfTheDaySorted.put(myObject.getTimestamp(), myObject))
       .peek(myObject -> lastObject.set(0, myObject))
       .anyMatch(predicate);

 objectsOfTheDaySorted.remove(lastObject.get(0).getTimestamp());

【讨论】:

  • 您可能希望使用set(0, myObject) 而不是add。否则,lastObject 列表会变得非常大。或者,将 lastObject 设为长度为 1 (lastObject[0] = myObject) 的数组。
  • @Malte Hartwig:感谢您的提示。尽管我实际上打算使用 set() 方法,但我完全忽略了这一点。我重新编辑了我的帖子。
  • 如果流从指定日期之前的对象开始,这将不起作用,因为谓词将匹配第一个并且已经终止流。此外,您使用 dayOfMonth,如果流包含超过 30 天的对象,则会导致错误。
  • 这是我以前知道的事实。我只关注上面的描述,它说整个流从一天开始,交付当天的所有对象,并在某个时候到达第二天。确实有优化空间。
  • @MalteHartwig 这是因为“This method 的存在主要是为了支持调试,您希望在其中看到元素流过管道中的某个点”。另见is peek really only for debugging?
【解决方案2】:

这取决于您是需要处理所有日子的对象还是某一天的对象。

基于 DiabolicWords 的回答,这是一个处理所有日子的示例:

TreeSet<MyObject> currentDaysObjects = new TreeSet<>(Comparator.comparing(MyObject::getTimestamp));
LocalDate[] currentDay = new LocalDate[1];
incoming.peek(o -> {
    LocalDate date = o.getTimestamp().toInstant().atZone(ZoneId.systemDefault()).toLocalDate();
    if (!date.equals(currentDay[0]))
    {
        if (currentDay != null)
        {
            processOneDaysObjects(currentDaysObjects);
            currentDaysObjects.clear();
        }
        currentDay[0] = date;
    }
}).forEach(currentDaysObjects::add);

这将收集一天的对象,处理它们,重置收集并在第二天继续。

如果您只想要某一天:

TreeSet<MyObject> currentDaysObjects = new TreeSet<>(Comparator.comparing(MyObject::getTimestamp));
LocalDate specificDay = LocalDate.now();
incoming.filter(o -> !o.getTimestamp()
                       .toInstant()
                       .atZone(ZoneId.systemDefault())
                       .toLocalDate()
                       .isBefore(specificDay))
        .peek(o -> currentDaysObjects.add(o))
        .anyMatch(o -> {
            if (o.getTimestamp().toInstant().atZone(ZoneId.systemDefault()).toLocalDate().isAfter(specificDay))
            {
                currentDaysObjects.remove(o);
                return true;
            }
            return false;
        });

过滤器将跳过specificDay之前的对象,而anyMatch将终止specificDay之后的流。

我已经读到在 Java 9 的流上会有类似 skipWhile 或 takeWhile 的方法。这些会让这变得容易得多。

在 Op 指定目标后编辑更详细

哇,这是一个很好的练习,而且很难破解。问题是一个明显的解决方案(收集流)总是贯穿整个流。您不能获取下一个 x 元素,对它们进行排序,对它们进行流式传输,然后在不一次对整个流(即所有天)执行此操作的情况下重复。出于同样的原因,在流上调用sorted() 将完全通过它(特别是当流不知道元素已经按天排序的事实时)。作为参考,请在此处阅读此评论:https://stackoverflow.com/a/27595803/7653073。

正如他们推荐的那样,这是一个包裹在流中的迭代器实现,它在原始流中向前看,获取一天的元素,对它们进行排序,然后在一个漂亮的新流中为您提供整个内容(不保留记忆中的所有日子!)。实现更加复杂,因为我们没有固定的块大小,但总是必须找到第二天的第一个元素才能知道何时停止。

public class DayByDayIterator implements Iterator<MyObject>
{
    private Iterator<MyObject> incoming;
    private MyObject next;

    private Iterator<MyObject> currentDay;

    private MyObject firstOfNextDay;
    private Set<MyObject> nextDaysObjects = new TreeSet<>(Comparator.comparing(MyObject::getTimestamp));

    public static Stream<MyObject> streamOf(Stream<MyObject> incoming)
    {
        Iterable<MyObject> iterable = () -> new DayByDayIterator(incoming);
        return StreamSupport.stream(iterable.spliterator(), false);
    }

    private DayByDayIterator(Stream<MyObject> stream)
    {
        this.incoming = stream.iterator();
        firstOfNextDay = incoming.next();
        nextDaysObjects.add(firstOfNextDay);
        next();
    }

    @Override
    public boolean hasNext()
    {
        return next != null;
    }

    @Override
    public MyObject next()
    {
        if (currentDay == null || !currentDay.hasNext() && incoming.hasNext())
        {
            nextDay();
        }

        MyObject result = next;

        if (currentDay != null && currentDay.hasNext())
        {
            this.next = currentDay.next();
        }
        else
        {
            this.next = null;
        }

        return result;
    }

    private void nextDay()
    {
        while (incoming.hasNext()
                && firstOfNextDay.getTimestamp().toLocalDate()
                .isEqual((firstOfNextDay = incoming.next()).getTimestamp().toLocalDate()))
        {
            nextDaysObjects.add(firstOfNextDay);
        }
        currentDay = nextDaysObjects.iterator();

        if (incoming.hasNext())
        {
            nextDaysObjects = new TreeSet<>(Comparator.comparing(MyObject::getTimestamp));
            nextDaysObjects.add(firstOfNextDay);
        }
    }
}

像这样使用它:

public static void main(String[] args)
{
    Stream<MyObject> stream = Stream.of(
            new MyObject(LocalDateTime.now().plusHours(1)),
            new MyObject(LocalDateTime.now()),
            new MyObject(LocalDateTime.now().plusDays(1).plusHours(2)),
            new MyObject(LocalDateTime.now().plusDays(1)),
            new MyObject(LocalDateTime.now().plusDays(1).plusHours(1)),
            new MyObject(LocalDateTime.now().plusDays(2)),
            new MyObject(LocalDateTime.now().plusDays(2).plusHours(1)));

    DayByDayIterator.streamOf(stream).forEach(System.out::println);
}

------------------- Output -----------------

2017-04-30T17:39:46.353
2017-04-30T18:39:46.333
2017-05-01T17:39:46.353
2017-05-01T18:39:46.353
2017-05-01T19:39:46.353
2017-05-02T17:39:46.353
2017-05-02T18:39:46.353

说明: currentDay 和 next 是迭代器的基础,而 firstOfNextDay 和 nextDaysObjects 已经查看了第二天的第一个元素。当currentDay 用尽时,调用nextDay() 并继续将incoming 的元素添加到nextDaysObjects 直到到达第二天,然后将nextDaysObjects 转换为currentDay。

一件事:如果传入的流是 null 或空的,它就会失败。你可以测试null,但是空的情况需要在工厂方法中捕获一个异常。为了便于阅读,我不想添加它。

我希望这是您需要的,请告诉我进展如何。

【讨论】:

  • 我的问题是我想在之后继续使用相同的流,我不太确定如何使用您建议的解决方案来实现。我已经更新了问题,以便更清楚地提及这一点。
  • @Johan 我添加了一个替代解决方案。它基本上使用 Iterator 而不是 Stream,但我添加了一个工厂 util 方法,该方法再次将该迭代器包装到一个流中。按如下方式使用:Stream&lt;MyObject&gt; sortedStream = DayByDayIterator.streamOf(unsortedStream)
【解决方案3】:

如果你考虑一种迭代方法,我认为它会变得更简单:

TreeSet<MyObject> currentDayObjects = new TreeSet<>(Comparator.comparing(MyObject::getTimestamp));
LocalDate currentDay = null;
for (MyObject m: stream::iterator) {
    LocalDate objectDay = m.getTimestamp().toInstant().atZone(ZoneId.systemDefault()).toLocalDate();
    if (currentDay == null) {
        currentDay = objectDay;
    } else if (!currentDay.equals(objectDay)) {
        // process a whole day of objects at once
        process(currentDayObjects);
        currentDay = objectDay;
        currentDayObjects.clear();
    }
    currentDayObjects.add(m);
}
// process the data of the last day
process(currentDayObjects);

【讨论】:

    猜你喜欢
    • 2016-08-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多