【问题标题】:Read from BufferedReader for a specific Duration从 BufferedReader 读取特定 Duration
【发布时间】:2023-03-18 20:05:02
【问题描述】:

所以,我正在从 BufferedReader 读取数据。一切都很好,直到我添加一个条件。我需要从 BufferedReader 读取特定的持续时间。

这就是我现在正在做的事情。

while ((line = br.readLine()) != null
                    && System.currentTimeMillis() - start < maxReadTime.toMillis()) { 
    // doingSomethingHere()
}

问题:即使在时间过去后,InputStream 仍然处于活动状态。 例如 - maxReadTime 是 30 秒。输入在 20 秒内不断出现。在接下来的 12 秒内,没有任何活动。现在,当下一个输入到达时,流打开并仅在读取输入后关闭。但是,我不处理此输入,因为 while 循环终止。

我的预期或我需要的是:Stream 将在 30 秒后关闭。也就是说,当输入到达第 32 秒时,流被关闭并且不监听任何输入。

我对 ExecutorService 知之甚少。我不确定这是否是正确的方法。

【问题讨论】:

  • 你需要先测试时间。否则,当计时器结束时,您仍然会阻止 readLine()。如果从套接字读取,您应该在套接字上设置读取超时而不是所有这些。
  • InputStream 可以来自任何东西。所以不是专门的套接字。那么我应该在 readLine() 之前运行另一个 while 循环吗?

标签: java inputstream bufferedreader


【解决方案1】:

只需在从流中读取之前设置您的计时器条件

while ((line = br.readLine()) != null) {
    boolean active = System.currentTimeMillis() - start < maxReadTime.toMillis();
    if (!active) {
        br.close();
    }         
    // doingSomethingHere()
}

在这种情况下,如果第一个条件是false(时间已到),则根本不会执行第二个条件

【讨论】:

  • 所以,InputStream 仍然处于活动状态。只有当我输入一些东西时它才会中断。我正在研究 Duration 结束时 InputStream 将关闭的地方。
  • 仍然无法正常工作。即使持续时间已过,我仍然可以接受输入。
【解决方案2】:

基本上,您必须在调用readLine() 之前通过调用方法ready() 检查缓冲区是否准备就绪,对于InputStream 检查available() 方法,该方法返回您可以在没有阻塞的情况下读取多少字节。

这里是一个例子

import java.io.*;
import java.time.Duration;

public class Main {

    public static void main(String[] args) {
        final InputStream in =  System.in; //new FileInputStream(new File("/tmp/x"));
        final String out = readInput(in, Duration.ofSeconds(5));
        System.out.printf("m=main, status=complete, out=%s%n", out);
    }

    public static String readInput(InputStream in, Duration duration) {
        final long timestamp = System.currentTimeMillis();
        final BufferedReader reader = new BufferedReader(new InputStreamReader(in));
        final StringBuilder out = new StringBuilder();
        try {
            String line = null;
            while (true){
                if(Duration.ofMillis(System.currentTimeMillis() - timestamp).compareTo(duration) >=0 ){
                    System.out.println("m=readInput, status=timeout");
                    break;
                }
                if(!reader.ready()){
                    System.out.println("m=readInput, status=not ready");
                    sleep(1000);
                    continue;
                }
                line = reader.readLine();
                if(line == null){
                    System.out.println("m=readInput, status=null line");
                    break;
                }
                out.append(line);
                out.append('\n');
                System.out.printf("m=readInput status=read, line=%s%n" , line);
            }
            return out.toString();
        } catch (IOException e){
            throw new RuntimeException(e);
        } finally {
            System.out.println("m=readInput, status=complete");
        }
    }

    static void sleep(int millis) {
        try {
            Thread.sleep(millis);
        } catch (InterruptedException e) {}
    }

}

如果您想在后台执行此操作,可以按照此示例进行操作

package com.mageddo;

import java.io.*;
import java.util.concurrent.*;

public class Main {

public static void main(String[] args) throws IOException, ExecutionException, InterruptedException {
        final InputStream in =  System.in; //new FileInputStream(new File("/tmp/x"));
        final StringBuilder out = new StringBuilder();
        final ExecutorService executor = Executors.newFixedThreadPool(1);
        final Future<String> promise = executor.submit(() -> readInput(in, out));
        try {
            final String result = promise.get(5, TimeUnit.SECONDS);
            System.out.printf("m=main, status=success, result=%s%n", result);
        } catch (TimeoutException e) {
            System.out.println("m=main, status=timeout");
            in.close();
            promise.cancel(true);
            System.out.println("Failed output: " + promise.get());
            e.printStackTrace();
        } finally {
            executor.shutdown();
            System.out.println("m=main, status=shutdown, out=" + out);
        }
    }

    public static String readInput(InputStream in, StringBuilder out) {
        final BufferedReader reader = new BufferedReader(new InputStreamReader(in));
        try {
            String line = null;
            while (true){
                if(Thread.currentThread().isInterrupted()){
                    System.out.println("m=readInput status=interrupt signal");
                    break;
                }
                if(!reader.ready()){
                    System.out.println("m=readInput, status=not ready");
                    sleep(1000);
                    continue;
                }
                line = reader.readLine();
                if(line == null){
                    System.out.println("m=readInput, status=null line");
                    break;
                }
                out.append(line);
                out.append('\n');
                System.out.printf("m=readInput status=read, line=%s%n" , line);
            }
            return out.toString();
        } catch (IOException e){
            throw new RuntimeException(e);
        } finally {
            System.out.println("m=readInput, status=complete");
        }
    }

    static void sleep(int millis) {
        try {
            Thread.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

}

See the reference

【讨论】:

  • 如果 InputStream 使用 @WIllNotClose 进行注释,这仍然有效吗?
  • 可以,因为解决方案跟close()方法无关,而是ready(),试试看
  • promise.get(5, TimeUnit.SECONDS) 是做什么的?
  • 它等待5秒然后如果超时则放弃,如果发生会抛出TimeoutException然后catch会将线程设置为中断然后readInput方法将停止从输入流读取跨度>
  • 那么我在哪里检查持续时间呢?在 readInput 方法中?
猜你喜欢
  • 1970-01-01
  • 2012-12-22
  • 1970-01-01
  • 2014-08-31
  • 2013-10-30
  • 2019-01-09
  • 2014-09-26
  • 1970-01-01
相关资源
最近更新 更多