【问题标题】:What is the (kind of) inverse operation to Java's Stream.flatMap()?Java 的 Stream.flatMap() 的(类型)逆操作是什么?
【发布时间】:2018-06-15 13:50:23
【问题描述】:

Stream.flatMap() 操作转换流

a, b, c

到每个输入元素包含零个或多个元素的流中,例如

a1, a2, c1, c2, c3

是否有相反的操作将几个元素组合成一个新元素?

  • 不是.reduce(),因为它只产生一个结果
  • 不是collect(),因为这只填充了一个容器(afau)
  • 它不是 forEach(),因为它只返回 void 并且有副作用

它存在吗?我可以用任何方式模拟它吗?

【问题讨论】:

  • 你在寻找什么返回类型?
  • Stream 进入,Stream 出现,其中 Y 是 X 的某种组合。原则上,整个事情与 collect() 非常相似,除了它会真正保持“流式传输”,而不是先收集然后流式传输结果:正如@Lino 的回答。
  • 你想要Collectors.groupingBy

标签: java java-stream


【解决方案1】:

最后我发现flatMap 可以说是它自己的“逆”。我监督flatMap 不一定会增加元素的数量。它还可以通过为某些元素发出空流来减少元素的数量。要实现分组操作,flatMap 调用的函数需要最小的内部状态,即最近的元素。它要么返回一个空流,要么在组结束时返回简化为组的代表。

这是一个快速实现,如果传入的两个元素不属于同一个组,则groupBorder 必须返回true,即它们之间是组边界。 combiner 是组合例如 (1,a)、(1,a)、(1,a) 到 (3,a) 的组函数,假设您的组元素是元组 (int, string) .

public class GroupBy<X> implements Function<X, Stream<X>>{

  private final BiPredicate<X, X> groupBorder;
  private final BinaryOperator<X> combiner;
  private X latest = null;

  public GroupBy(BiPredicate <X, X> groupBorder,
                 BinaryOperator<X> combiner) {
    this.groupBorder = groupBorder;
    this.combiner = combiner;
  }

  @Override
  public Stream<X> apply(X elem) {
    // TODO: add test on end marker as additonal parameter for constructor
    if (elem==null) {
      return latest==null ? Stream.empty() : Stream.of(latest);
    }
    if (latest==null) {
      latest = elem;
      return Stream.empty();
    }
    if (groupBorder.test(latest, elem)) {
      Stream<X> result = Stream.of(latest);
      latest = elem;
      return result;
    }
    latest = combiner.apply(latest,  elem);
    return Stream.empty();
  }
}

但有一个警告:要发送整个流的最后一组,必须将结束标记作为最后一个元素粘贴到流中。上面的代码假定它是null,但可以添加一个额外的结束标记测试器。

我想不出不依赖于结束标记的解决方案。

此外,我也没有在传入和传出元素之间进行转换。对于唯一操作,这将起作用。对于计数操作,上一步必须将单个元素映射到计数对象。

【讨论】:

  • 这个功能的实际用途是什么?
【解决方案2】:

你可以破解你的方式。请参见以下示例:

Stream<List<String>> stream = Stream.of("Cat", "Dog", "Whale", "Mouse")
   .collect(Collectors.collectingAndThen(
       Collectors.partitioningBy(a -> a.length() > 3),
       map -> Stream.of(map.get(true), map.get(false))
    ));

【讨论】:

  • :-) 确实是个黑客。
【解决方案3】:
    IntStream.range(0, 10)
            .mapToObj(n -> IntStream.of(n, n / 2, n / 3))
            .reduce(IntStream.empty(), IntStream::concat)
            .forEach(System.out::println);

如您所见,元素也映射到 Streams,然后连接成一个大流。

【讨论】:

  • 我可能对此有误解或无法进行传输,但这似乎假设我在看到a1时可以生成a2、a3和a4。但这种情况并非如此。出现了一些任意元素,我想将它们批量处理。
  • @Harald 我理解:编辑几个流式元素的子范围,改变流。就像管道 I/O / ... 中的数据过滤器一样。我认为 Stream 不适合。
【解决方案4】:

看看collapse中的StreamEx

StreamEx.of("a1", "a2", "c1", "c2", "c3").collapse((a, b) -> a.charAt(0) == b.charAt(0))
    .map(e -> e.substring(0, 1)).forEach(System.out::println);

或者my fork具有更多功能:groupBysplitsliding...

StreamEx.of("a1", "a2", "c1", "c2", "c3").collapse((a, b) -> a.charAt(0) == b.charAt(0))
.map(e -> e.substring(0, 1)).forEach(System.out::println);
// a
// c

StreamEx.of("a1", "a2", "c1", "c2", "c3").splitToList(2).forEach(System.out::println);
// [a1, a2]
// [c1, c2]
// [c3]

StreamEx.of("a1", "a2", "c1", "c2", "c3").groupBy(e -> e.charAt(0))
.forEach(System.out::println);
// a=[a1, a2]
// c=[c1, c2, c3]

【讨论】:

    【解决方案5】:

    这是我想出的:

    interface OptionalBinaryOperator<T> extends BiFunction<T, T, Optional<T>> {
      static <T> OptionalBinaryOperator<T> of(BinaryOperator<T> binaryOperator,
              BiPredicate<T, T> biPredicate) {
        return (t1, t2) -> biPredicate.test(t1, t2)
                ? Optional.of(binaryOperator.apply(t1, t2))
                : Optional.empty();
      }
    }
    
    class StreamUtils {
      public static <T> Stream<T> reducePartially(Stream<T> stream,
              OptionalBinaryOperator<T> conditionalAccumulator) {
        Stream.Builder<T> builder = Stream.builder();
        stream.reduce((t1, t2) -> conditionalAccumulator.apply(t1, t2).orElseGet(() -> {
          builder.add(t1);
          return t2;
        })).ifPresent(builder::add);
        return builder.build();
      }
    }
    

    不幸的是,我没有时间让它变得懒惰,但可以通过编写一个自定义的Spliterator 委托给stream.spliterator() 来完成,这将遵循上述逻辑(而不是使用stream.reduce(),这是一个终端操作)。


    PS。我刚刚意识到你想要&lt;T,U&gt; 转换,我写了关于&lt;T,T&gt; 转换的文章。如果你可以先从T映射到U,然后使用上面的函数,那就这样了(即使是次优的)。

    如果它更复杂,则需要在提出 API 之前定义减少/合并的条件(例如 Predicate&lt;T&gt;BiPredicate&lt;T,T&gt;BiPredicate&lt;U,T&gt;,甚至可能是 Predicate&lt;List&lt;T&gt;&gt;)。

    【讨论】:

      【解决方案6】:

      有点像 StreamEx,你可以手动实现 Spliterator。例如,

      collectByTwos(Stream.of(1, 2, 3, 4), (x, y) -> String.format("%d%d", x, y))
      

      ...使用以下代码返回“12”、“34”流:

      public static <X,Y> Stream<Y> collectByTwos(Stream<X> inStream, BiFunction<X,X,Y> mapping) {
          Spliterator<X> origSpliterator = inStream.spliterator();
          Iterator<X> origIterator = Spliterators.iterator(origSpliterator);
      
          boolean isParallel = inStream.isParallel();
          long newSizeEst = (origSpliterator.estimateSize() + 1) / 2;
      
          Spliterators.AbstractSpliterator<Y> lCombinedSpliterator =
                  new Spliterators.AbstractSpliterator<>(newSizeEst, origSpliterator.characteristics()) {
              @Override
              public boolean tryAdvance(Consumer<? super Y> action) {
                  if (! origIterator.hasNext()) {
                      return false;
                  }
                  X lNext1 = origIterator.next();
                  if (! origIterator.hasNext()) {
                      throw new IllegalArgumentException("Trailing elements of the stream would be ignored.");
                  }
                  X lNext2 = origIterator.next();
                  action.accept(mapping.apply(lNext1, lNext2));
                  return true;
              }
          };
          return StreamSupport.stream(lCombinedSpliterator, isParallel)
                  .onClose(inStream::close);
      }
      

      (我认为这对于并行流来说可能是不正确的。)

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2015-03-08
        • 2015-05-20
        • 1970-01-01
        • 2015-05-28
        • 1970-01-01
        • 1970-01-01
        • 2011-06-04
        相关资源
        最近更新 更多