【问题标题】:Aggregation in JAVA with Streams, is this a good approach?在 JAVA 中使用 Streams 进行聚合,这是一个好方法吗?
【发布时间】:2017-10-17 08:35:42
【问题描述】:

我正在玩一点 Java 流,我想出了一个解决问题的方法,我想与你分享,看看我的方法是否正确。

我从https://catalog.data.gov/dataset/consumer-complaint-database 下载了一个数据集,其中包含超过 70 万条客户投诉记录。 我使用的信息如下:

公司名称 产品名称

我的目标是获得以下结果:

数据集中出现次数较多的 10 家公司

数据集中出现次数较多的 10 个产品

得到类似的东西

Map<String, Map<String,Integer>>

其中,主地图的key为公司名称,二级地图的key为产品名称,其值为产品在该公司被投诉的次数。

所以我所做的解决方案如下:

@Test
public void joinGroupingsTest() throws URISyntaxException, IOException {
    String path = CsvReaderTest.class.getResource("/complains.csv").toURI().toString();
    complains = CsvReader.readFileStreamComplain(path.substring(path.indexOf('/')+1));

    Map<String, List<Complain>> byCompany = complains.parallelStream()
            .collect(Collectors.groupingBy(Complain::getCompany))
            .entrySet().stream()
            .sorted((f1, f2) -> Long.compare(f2.getValue().size(), f1.getValue().size()))
            .limit(10)
            .collect(Collectors.toMap(Entry::getKey, Entry::getValue));


    Map<String, List<Complain>> byProduct = complains.parallelStream()
            .collect(Collectors.groupingBy(Complain::getProduct))
            .entrySet().stream()
            .sorted((f1, f2) -> Long.compare(f2.getValue().size(), f1.getValue().size()))
            .limit(10)
            .collect(Collectors.toMap(Entry::getKey, Entry::getValue));

    Map<String, List<Complain>> map = complains.parallelStream()
            .filter((x) -> byCompany.get(x.getCompany()) != null
                && byProduct.get(x.getProduct()) != null)
            .collect(Collectors.groupingBy(Complain::getCompany));

    Map<String, Map<String, Long>> map2 = map.entrySet().parallelStream()
            .collect(Collectors.toMap(
                    e -> e.getKey(),
                    e -> e.getValue().stream()
                            .collect(Collectors.groupingBy(Complain::getProduct, Collectors.counting()))
            ));


   System.out.println(map2);


}

如您所见,我有几个步骤可以实现这一目标:

1) 我得到了 10 家发生次数较多的公司以及相关的投诉(记录)

2) 我得到了 10 种出现次数较多的产品以及相关的投诉(记录)

3) 我得到一张以公司名称为关键字的地图,该地图在之前计算的前 10 名公司中以及同样在前 10 名产品中的产品的抱怨

4) 我进行所需的转换以获得我想要的地图。

除了在两个不同的线程中分叉和分离步骤 1 和 2 之外,是否还有其他考虑因素需要我提高性能甚至以更好的方式使用流。

谢谢!

【问题讨论】:

  • 如果数据集中投诉较多的 10 种产品没有被投诉较多的 10 家公司中的任何一家销售怎么办?我认为这是您设计解决方案时的错误。此外,您正在遍历整个数据集 3 次,这看起来不是很理想,恕我直言
  • 使用流进行聚合并没有什么坏处,但是,使用parallelStream 可能会产生开销。我经常在我的 API 中使用stream。这是一个相关的线程:stackoverflow.com/questions/20375176/…
  • 没有理由分开最后两个流操作,只需将最后一个的groupingBy收集器作为第二个参数传递给前一个的groupingBy收集器。
  • @Holger 你是对的,这肯定是一个修复,我不知道我能做到这一点

标签: java performance java-8 java-stream collectors


【解决方案1】:

在前两个操作中,您将集合组放入Lists,只是为了按它们的大小排序。这显然是一种资源浪费,因为您可以在分组时简单地计算组元素,然后按计数排序。此外,由于前两个操作是相同的,除了分组功能外,值得通过为任务创建方法来消除代码重复。

其他两个流操作可以合二为一,通过在收集组时立即对组执行collect 操作。

public void joinGroupingsTest() throws URISyntaxException, IOException {
    String path = CsvReaderTest.class.getResource("/complains.csv").toURI().toString();
    complains = CsvReader.readFileStreamComplain(path.substring(path.indexOf('/')+1));

    Set<String> byCompany = getTopTen(complains, Complain::getCompany);
    Set<String> byProduct = getTopTen(complains, Complain::getProduct);

    Map<String, Map<String, Long>> map = complains.stream()
            .filter(x -> byCompany.contains(x.getCompany())
                      && byProduct.contains(x.getProduct()))
            .collect(Collectors.groupingBy(Complain::getCompany,
                Collectors.groupingBy(Complain::getProduct, Collectors.counting())));
   System.out.println(map);
}
static <T,V> Set<V> getTopTen(Collection<T> source, Function<T,V> criteria) {
    return source.stream()
            .collect(Collectors.groupingBy(criteria, Collectors.counting()))
            .entrySet().stream()
            .sorted(Map.Entry.comparingByValue())
            .limit(10)
            .map(Map.Entry::getKey)
            .collect(Collectors.toSet());
}

请注意,这两个条件的交集可能小于十个元素,甚至可能为空。您可能会重新考虑条件。

此外,您应该经常重新检查数据量是否真的足够大以受益于并行处理。另请注意,getTopTen 操作由两个流操作组成。将第一个切换为并行不会改变第二个的性质。

【讨论】:

  • 我觉得我为此做了太多的操作。
  • 非常感谢您的回答,我将阅读更多有关流的信息,以尝试以这种方法进行思考。你知道一本 100% 致力于此的好书吗?
【解决方案2】:

如果您不必处理性能问题,使用 java 流是一个很好的方法。 Java Stream 或并行流相对较慢,如果在流中进行操作时出现任何异常,可能会使您的调试变得最糟糕。流的好处是您必须编写几行代码来解决复杂的聚合问题或数据结构更改。这是一个链接,您可以在其中了解 java 流与传统方法相比有多慢。

https://blog.codefx.org/java/stream-performance/

【讨论】:

猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-11-05
  • 1970-01-01
  • 2012-07-20
  • 2018-07-18
  • 2021-01-16
相关资源
最近更新 更多