小编Woj*_*ela的帖子

RxJava:在onNext中调用取消订阅

我想知道unsubscribe从这样的onNext处理程序中调用是否合法:

List<Integer> gatheredItems = new ArrayList<>();

Subscriber<Integer> subscriber = new Subscriber<Integer>() {
    public void onNext(Integer item) {
        gatheredItems.add(item);
        if (item == 3) {
            unsubscribe();
        }
    }
    public void onCompleted() {
        // noop
    }
    public void onError(Throwable sourceError) {
        // noop
    }
};

Observable<Integer> source = Observable.range(0,100);

source.subscribe(subscriber);
sleep(1000);
System.out.println(gatheredItems);
Run Code Online (Sandbox Code Playgroud)

上面的代码正确输出只收集了四个元素:[0, 1, 2, 3].但是如果有人将源observable更改为缓存:

Observable<Integer> source = Observable.range(0,100).cache();
Run Code Online (Sandbox Code Playgroud)

然后收集所有一百个元素.我没有对source observable的控制(无论是否缓存),那么如何从内部取消订阅onNext呢?

顺便说一句:那么在不onNext正确的事情中取消订阅是做什么的?

(我的实际用例是,onNext我实际上是在写输出流,当IOException发生时,没有什么可以写入输出,所以我需要以某种方式停止进一步处理.)

java rx-java

10
推荐指数
1
解决办法
2752
查看次数

在HTTP Servlet中正确流送输入和输出

我正在尝试编写将处理POST请求并流式传输输入和输出的servlet。我的意思是它应该读取一行输入,在该行上做一些工作,然后写入一行输出。并且它应该能够处理任意长请求(这样也会产生任意长响应)而不会出现内存不足异常。这是我的第一次尝试:

protected void doPost(HttpServletRequest request, HttpServletResponse response) {
    ServletInputStream input = request.getInputStream();
    ServletOutputStream output = response.getOutputStream();

    LineIterator lineIt = lineIterator(input, "UTF-8");
    while (lineIt.hasNext()) {
        String line = lineIt.next();
        output.println(line.length());
    }
    output.flush();
}
Run Code Online (Sandbox Code Playgroud)

现在,我使用来测试了该servlet curl,它可以工作,但是当我使用Apache HttpClient编写客户端时,客户端线程和服务器线程都会挂起。客户看起来像这样:

HttpClient client = HttpClientBuilder.create().build();
HttpPost post = new HttpPost(...);

// request
post.setEntity(new FileEntity(new File("some-huge-file.txt")));
HttpResponse response = client.execute(post);

// response
copyInputStreamToFile(response.getEntity().getContent(), new File("results.txt"));
Run Code Online (Sandbox Code Playgroud)

问题很明显。客户端在一个线程中按顺序执行它的工作-首先它完全发送请求,然后才开始读取响应。但是服务器为每行输入写入一行输出,如果客户端未读取输出(而顺序客户端未读取),则服务器将被阻止尝试写入输出流。反过来,这会阻止客户端尝试将输入发送到服务器。

我猜想是curl有效的,因为它以某种方式同时发送输入和接收输出(在单独的线程中?)。因此,第一个问题是是否可以将Apache HttpClient配置为与以下行为类似curl?

下一个问题是,如何改进servlet,以便使行为不佳的客户端不会导致服务器线程挂起?我的第一个尝试是引入中间缓冲区,该缓冲区将收集输出,直到客户端完成发送输入为止,然后servlet才开始发送输出:

ServletInputStream input = request.getInputStream();
ServletOutputStream output = response.getOutputStream();

// prepare intermediate store
int …
Run Code Online (Sandbox Code Playgroud)

java streaming multithreading servlets http

5
推荐指数
1
解决办法
1242
查看次数

RxJava:如何使用我自己的边界函数对任意大量的项进行分组

我有一个可以发出字符串的observable,我想用第一个字符对它们进行分组.这样做很容易groupBy:

Observable<String> rows = Observable.just("aa", "ab", "ac", "bb", "bc", "cc");

Observable<List<String>> groupedRows = rows.groupBy(new Func1<String, Character>() {
  public Character call(String row) {
    return row.charAt(0);
  }
}).flatMap(new Func1<GroupedObservable<Character, String>, Observable<List<String>>>() {
  public Observable<List<String>> call(GroupedObservable<Character, String> group) {
    return group.toList();
  }
});

groupedRows.toBlocking().forEach(new Action1<List<String>>() {
  public void call(List<String> group) {
    System.out.println(group);
  }
});

// Output:
// [aa, ab, ac]
// [bb, bc]
// [cc]
Run Code Online (Sandbox Code Playgroud)

但它对我的目的并不好,因为groupBy只有在源可观察量发出时才完成每个组onComplete.因此,如果我有很多行,它们将完全聚集在内存中,并且只在最后一行"刷新"并写入输出.

我需要像buffer运算符这样的东西,但是我自己的函数表示每个组的边界.我实现了它(知道行总是按字母排序):

Observable<String> rows = Observable.just("aa", "ab", "ac", "bb", "bc", …
Run Code Online (Sandbox Code Playgroud)

java rx-java

5
推荐指数
1
解决办法
796
查看次数

标签 统计

java ×3

rx-java ×2

http ×1

multithreading ×1

servlets ×1

streaming ×1