【问题标题】:Java 8 streams - incremental collection/partial reduce/intermittent map/...what is this even called?Java 8 流 - 增量收集/部分减少/间歇映射/......这甚至叫什么?
【发布时间】:2015-02-01 16:59:22
【问题描述】:

我正在处理遵循该模式的可能无限的数据元素流:

E1 <start mark> E2 foo E3 bah ... En-1 bar En <end mark>

也就是说,一个 的流,必须先在缓冲区中累积,然后才能将它们映射到对象模型。

目标:将Stream&lt;String&gt; 聚合成Stream&lt;ObjectDefinedByStrings&gt;没有在无限流上收集的开销。

在英语中,代码类似于“一旦你看到一个开始标记,开始缓冲。缓冲直到你看到一个结束标记,然后准备返回旧缓冲区,并准备一个新缓冲区。返回旧缓冲区。”

我当前的实现形式为:

Data<String>.stream()
            .map(functionReturningAnOptionalPresentOnlyIfObjectIsComplete)
            .filter(Optional::isPresent)

我有几个问题:

  1. 这个操作的正确名称是什么? (即,我可以用 Google 搜索更多示例吗?我发现的每个关于 .map() 的讨论都谈到了 1:1 映射。每个关于 .reduce 的讨论)都谈到了 n:1 缩减。 .collect()的每一次讨论都谈到了积累作为终端操作...)

  2. 这在很多方面看起来都很糟糕。有没有更好的方法来实现这一点? (.collectUntilConditionThenApplyFinisher(Collector,Condition,Finisher)...形式的候选人?)

谢谢!

【问题讨论】:

  • 这似乎是一个非常糟糕的主意。 map 操作应该没有副作用。您几乎可以肯定应该做的是调用Stream.iterator(),并通过移动迭代器以“老派”方式执行此操作,直到您点击&lt;end&gt;
  • 或多或少,流并不打算以这种方式使用。迭代器是一种更合理的方式。
  • facepalm 这完全正确。 #FluCoding

标签: java dictionary java-8 java-stream reduce


【解决方案1】:

为避免混淆,您可以在映射之前进行过滤。

Data<String>.stream()
    .filter(text -> canBeConvertedToObject(text))
    .map(text -> convertToObject(text))

这在无限流上运行得非常好,并且只构造需要构造的对象。它还避免了创建不必要的 Optional 对象的开销。

【讨论】:

    【解决方案2】:

    不幸的是,Java 8 Stream API 中没有部分归约操作。然而,这样的操作是在我的StreamEx 库中实现的,它增强了标准的 Java 8 Streams。所以你的任务可以这样解决:

    Stream<ObjectDefinedByStrings> result = 
        StreamEx.of(strings)
                .groupRuns((a, b) -> !b.contains("<start mark>"))
                .map(stringList -> constructObjectDefinedByStrings());
    

    strings 是普通的 Java-8 流或其他源,如数组、CollectionSpliterator 等。适用于无限或并行流。 groupRuns 方法采用 BiPredicate,它应用于两个相邻的流元素,如果必须对这些元素进行分组,则返回 true。这里我们说元素应该被分组,除非第二个包含"&lt;start mark&gt;"(这是新元素的开始)。之后,您将获得List&lt;String&gt; 元素的流。

    如果收集到中间列表不适合您,您可以使用collapse(BiPredicate, Collector) 方法并指定自定义收集器来执行部分缩减。例如,您可能希望将所有字符串连接在一起:

    Stream<ObjectDefinedByStrings> result = 
        StreamEx.of(strings)
                .collapse((a, b) -> !b.contains("<start mark>"), Collectors.joining())
                .map(joinedString -> constructObjectDefinedByStrings());
    

    【讨论】:

      【解决方案3】:

      我为这种部分减少提出了另外 2 个用例:

      1。解析 SQL 和 PL/SQL(Oracle 过程)语句

      SQL 语句的标准分隔符是分号 (;)。它将普通的 SQL 语句彼此分开。但是如果你有 PL/SQL 语句,那么分号将语句中的运算符相互分隔,而不仅仅是整个语句。

      解析包含普通 SQL 和 PL/SQL 语句的脚本文件的一种方法是首先用分号将它们拆分,然后如果特定语句以特定关键字(DECLAREBEGIN 等)开头,则加入此带有遵循 PL/SQL 语法规则的 next 语句的语句。

      顺便说一句,这不能通过使用StreamEx 部分归约操作来完成,因为它们只测试两个相邻的元素。由于您需要了解从初始 PL/SQL 关键字元素开始的先前流元素,以确定是否将当前元素包含到部分归约中或应该完成部分归约。在这种情况下,可变部分归约可用于收集器持有已收集元素的信息,并且某些Predicate 仅测试收集器本身(如果应该完成部分归约)或BiPredicate 测试收集器和当前流元素。

      理论上,我们谈论的是使用流管道意识形态实现 LR(0) 或 LR(1) 解析器(请参阅https://en.wikipedia.org/wiki/LR_parser)。 LR-parser 可用于解析大多数编程语言的语法。

      解析器是一个带栈的有限自动机。在 LR(0) 自动机的情况下,它的转换仅取决于堆栈。在 LR(1) 自动机的情况下,它取决于堆栈和流中的下一个元素(理论上可以有 LR(2)、LR(3) 等。自动机偷看 2、3 等下一个元素来确定转换,但实际上,所有编程语言在语法上都是 LR(1) 语言)。

      要实现解析器,应该有一个 Collector 包含有限自动机的堆栈和谓词测试是否达到了该自动机的最终状态(这样我们就可以停止归约)。在 LR(0) 的情况下,它应该是 Predicate 测试 Collector 本身。在 LR(1) 的情况下,它应该是 BiPredicate 测试 Collector 和流中的下一个元素(因为转换取决于堆栈和下一个符号)。

      所以要实现 LR(0) 解析器,我们需要类似以下内容(T 是流元素类型,A 是同时保存有限自动机堆栈和结果的累加器,R 是每个解析器工作形成输出的结果流):

      <R,A> Stream<R> Stream<T>.parse(
          Collector<T,A,R> automataCollector,
          Predicate<A> isFinalState)
      

      (为了紧凑,我删除了 ? super T 之类的复杂性,而不是 T - 结果 API 应该包含这些)

      要实现 LR(1) 解析器,我们需要以下内容:

      <R,A> Stream<R> Stream<T>.parse(
          BiPredicate<A, T> isFinalState
          Collector<T,A,R> automataCollector)
      

      注意:在这种情况下,BiPredicate 应该元素被累加器消耗之前测试它。请记住 LR(1) 解析器正在查看下一个元素以确定转换。因此,如果空累加器拒绝接受下一个元素(BiPredicate 返回 true,表示部分归约结束,在刚刚由 Supplier 和下一个流元素创建的空累加器上),则可能存在潜在异常。

      2。基于流元素类型的条件批处理

      当我们执行 SQL 语句时,我们希望将相邻的数据修改 (DML) 语句合并到一个批处理中(请参阅 JDBC API)以提高整体性能。但是我们不想批量查询。所以我们需要条件批处理(而不是像Java 8 Stream with batch processing 那样的无条件批处理)。

      对于这种特定情况,可以使用StreamEx 部分归约操作,因为如果BiPredicate 测试的两个相邻元素都是 DML 语句,则它们应该包含在批处理中。所以我们不需要知道以前的批次收集历史。

      但是我们可以增加任务的复杂性,并说批次应该受到大小的限制。比如说,一批不超过 100 个 DML 语句。在这种情况下,我们不能忽略以前的批处理收集历史记录,并且使用BiPredicate 来确定应该继续还是停止批处理收集是不够的。

      虽然我们可以在 StreamEx 部分归约之后添加 flatMap 以将长批次拆分为多个部分。但这会延迟特定的 100 元素批处理执行,直到所有 DML 语句都被收集到无限批处理中。毋庸置疑,这违反了管道意识形态:我们希望最小化缓冲以最大化输入和输出之间的速度。此外,如果 DML 语句的列表非常长且中间没有任何查询(例如,数据库导出导致数百万个 INSERTs),无限批量收集可能会导致 OutOfMemoryError,这是无法容忍的。

      因此,在这种具有上限的复杂条件批处理收集的情况下,我们还需要像前面用例中描述的 LR(0) 解析器一样强大的东西。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-02-17
        • 2015-02-24
        • 2017-09-22
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多