【问题标题】:RxJava buffer/window until element suffers conditionRxJava 缓冲区/窗口,直到元素遭受条件
【发布时间】:2015-07-09 10:06:28
【问题描述】:

我正在尝试从阅读器中读取行并将它们分组到属于一起的块。

原文:

bla1
bla2
### block separator ###
bla3
bla4
### block separator ###
...

我需要得到两个块(bla1,bla2)和(bla3,bla4)。

代码:

import org.apache.commons.lang3.StringUtils;
import rx.Observable;

import java.io.BufferedReader;
import java.io.FileReader;
import java.io.Reader;
import java.util.Iterator;

public class BlockBuilder {

  public static void main(String[] args) {

    try {
      FileReader fileReader = new FileReader("/path/to/some/file");
      LineIterable lineIterable = new LineIterable(fileReader);

      Observable.from(lineIterable)
          .buffer(100)
          // Needed instead of time/count: until line matches condition
          // something like .buffer(line -> line.equals("### block separator ###")
          .forEach(gatheredLines -> {
            String gatheredBlock = StringUtils.join(gatheredLines, '\n');
            System.out.println(gatheredBlock);
            System.out.println("###### ###### ###### ######");
          });
    } catch (Exception ex) {
      ex.printStackTrace();
    }
  }

  private static class LineIterable implements Iterable<String> {
    private final Iterator<String> iterator;
    public LineIterable(Reader reader) {
      iterator = new BufferedReader(reader).lines().iterator();
    }
    @Override
    public Iterator<String> iterator() {
      return iterator;
    }
  } 
}

使用缓冲区或窗口无关紧要,或者我对这两者的看法完全错误。

我认为缓冲区的 bufferClosingSelector 或窗口的 closingSelector 一定是可能的。 两者都是创建 Observer 的函数,它可以触发关闭当前缓冲区或窗口,但我看不到在哪里可以获取当前行。

【问题讨论】:

    标签: rx-java


    【解决方案1】:

    您可以发布您的源代码并将其用于缓冲和缓冲区边界:

    Observable<String> source = Observable.just(
            "a", "b", "#", 
            "c", "d", "e", "#", 
            "f", "g");
    
    source.publish(p -> 
            p.filter(v -> !"#".equals(v))
            .buffer(() -> p.filter(v -> "#".equals(v))))
    .subscribe(System.out::println);
    

    【讨论】:

    • 那太好了,它编译时没有错误。但是,一旦我尝试此代码,我就会收到一个巨大的错误: Error:(16, 25) java: no suitable method found for buffer(()->p.filt[...]s(v))) method rx .Observable.buffer(rx.functions.Func0 extends rx.Observable extends TClosing>>) 不适用(无法推断类型变量 TClosing(参数不匹配;lambda 表达式 rx 中的返回类型错误.Observable 无法转换为 ? extends rx.Observable extends TClosing>)) [...]
    • 从 Eclipse 4.5 开始为我工作。您可能需要在第二个 p.filter() 调用之前添加显式 &lt;String&gt;
    • 有人能解释一下publish()在这种情况下的用途吗?
    • 源的信号必须组播到多个子流,不能多次订阅源(不保证发出相同的数据)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-10-28
    • 2020-02-26
    • 1970-01-01
    • 2023-03-07
    相关资源
    最近更新 更多