【问题标题】:submit method won't invoke onNext FLOW STREAM API JAVA提交方法不会调用 onNext FLOW STREAM API JAVA
【发布时间】:2019-10-01 22:57:57
【问题描述】:

我正在学习 Java 中的 FLOW Stream API,我目前正在创建一个基于 oracle community 的示例。问题是我看不到预期的输出,而只是看到onSubscribe 方法中打印的订阅字符串。我已经在 StackOverflow 上检查并找到了 submissionpublisher-on-submit-not-invoking-onnext-of-subscriber,但没有工作,因为我已经在调用 request(Long N)

import java.util.concurrent.Flow;

public class Computer<T> implements Flow.Subscriber<T> {

    private Flow.Subscription subscription;

    @Override
    public void onSubscribe(Flow.Subscription subscription) {
        this.subscription = subscription;
        System.out.println("SUBSCRIBING");
        this.subscription.request(1);
    }

    @Override
    public void onNext(T item) {
        System.out.println(String.format("Got %s", item.toString()));
        this.subscription.request(1);
    }

    @Override
    public void onError(Throwable throwable) {
        throwable.printStackTrace();
    }

    @Override
    public void onComplete() {
        System.out.println("DONE");
    }

}

--

import java.util.List;
import java.util.concurrent.SubmissionPublisher;

public class Sensor {

    public static void main(String[] args) {
        SubmissionPublisher<String> submissionPublisher = new SubmissionPublisher<>();
        Computer<String> subscriber = new Computer<>();
        submissionPublisher.subscribe(subscriber);

        List<String> items = List.of("1.25", "1.224", "1.55");
        items.forEach(submissionPublisher::submit);
        submissionPublisher.close();
    }

}

我只是去看看:

SUBSCRIBING

为什么onNext 方法没有被调用?

【问题讨论】:

    标签: java reactive-programming publish-subscribe java-9


    【解决方案1】:

    您不会将ScheduledExecutorService 传递给Publisher,它基本上是一个 ExecutorService,它可以安排任务在延迟后运行或在每次执行之间以固定的时间间隔重复执行。

    import java.util.List;
    import java.util.concurrent.Executors;
    import java.util.concurrent.ScheduledExecutorService;
    import java.util.concurrent.SubmissionPublisher;
    
    public class Sensor {
    
        public static void main(String[] args) {
            ScheduledExecutorService executor = Executors.newScheduledThreadPool(Runtime.getRuntime().availableProcessors());
            SubmissionPublisher<String> submissionPublisher = new SubmissionPublisher<>(executor, 5);
            Computer<String> subscriber = new Computer<>();
            submissionPublisher.subscribe(subscriber);
    
            List<String> items = List.of("1.25", "1.224", "1.55");
            items.forEach(submissionPublisher::submit);
            submissionPublisher.close();
            executor.shutdown();
        }
    }
    

    【讨论】:

    • 如何配置传递的执行者,让发布者以固定速率发布?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-01-13
    • 2020-07-08
    • 2013-10-01
    相关资源
    最近更新 更多