【问题标题】:Java 9 Behavior of Flow SubmissionPublisher offer methodFlow SubmissionPublisher 提供方法的 Java 9 行为
【发布时间】:2017-09-30 11:15:42
【问题描述】:

我一直在玩 Java Flow offer 运算符,但在阅读文档并进行测试后我不明白。

这是我的测试

@Test
public void offer() throws InterruptedException {
    //Create Publisher for expected items Strings
    SubmissionPublisher<String> publisher = new SubmissionPublisher<>();
    //Register Subscriber
    publisher.subscribe(new CustomSubscriber<>());
    publisher.subscribe(new CustomSubscriber<>());
    publisher.subscribe(new CustomSubscriber<>());
    publisher.offer("item", (subscriber, value) -> false);
    Thread.sleep(500);
}

offer 运算符接收要发出的项目和 BiPredicate 函数,据我了解阅读文档,只有在谓词函数为真的情况下才会发出项目。

bur通过测试后的结果是

Subscription done:
Subscription done:
Subscription done:
Got : item --> onNext() callback
Got : item --> onNext() callback
Got : item --> onNext() callback

如果我返回 true 而不是 false,则结果没有变化。

请谁能给我解释一下这个运算符。

【问题讨论】:

    标签: java java-9 reactive-streams java-flow


    【解决方案1】:

    不,谓词函数用于决定是否重试docs中提到的发布操作:

    onDrop - 如果非空,则在拖放到订阅者时调用的处理程序,带有订阅者和项目的参数;如果返回 true,则重新尝试(一次)提议

    它不影响最初是否发送该项目。

    编辑:使用offer 方法时如何发生掉落的示例

    我想出了一个示例,说明在调用 offer 方法时如何发生丢包。我不认为输出是 100% 确定的,但是当它运行几次时有明显的区别。您可以将处理程序更改为返回 true 而不是 false,以查看重试如何减少由于饱和缓冲区造成的丢弃。在此示例中,通常会发生下降,因为最大缓冲区容量明显很小(传递给 SubmissionPublisher 的构造函数)。但是,当在一小段睡眠期后启用重试时,会删除丢弃:

    public class SubmissionPubliserDropTest {
    
        public static void main(String[] args) throws InterruptedException {
            // Create Publisher for expected items Strings
            // Note the small buffer max capacity to be able to cause drops
            SubmissionPublisher<String> publisher =
                                   new SubmissionPublisher<>(ForkJoinPool.commonPool(), 2);
            // Register Subscriber
            publisher.subscribe(new CustomSubscriber<>());
            publisher.subscribe(new CustomSubscriber<>());
            publisher.subscribe(new CustomSubscriber<>());
            // publish 3 items for each subscriber
            for(int i = 0; i < 3; i++) {
                int result = publisher.offer("item" + i, (subscriber, value) -> {
                    // sleep for a small period before deciding whether to retry or not
                    try {
                        Thread.sleep(200);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                    return false;  // you can switch to true to see that drops are reduced
                });
                // show the number of dropped items
                if(result < 0) {
                    System.err.println("dropped: " + result);
                }
            }
            Thread.sleep(3000);
            publisher.close();
        }
    }
    
    class CustomSubscriber<T> implements Flow.Subscriber<T> {
    
        private Subscription sub;
    
        @Override
        public void onComplete() {
            System.out.println("onComplete");
        }
    
        @Override
        public void onError(Throwable th) {
            th.printStackTrace();
            sub.cancel();
        }
    
        @Override
        public void onNext(T arg0) {
            System.out.println("Got : " + arg0 + " --> onNext() callback");
            sub.request(1);
        }
    
        @Override
        public void onSubscribe(Subscription sub) {
            System.out.println("Subscription done");
            this.sub = sub;
            sub.request(1);
        }
    
    }
    

    【讨论】:

    • 我可能很笨,但是在修改我的测试并再次阅读之后,文档仍然不明白。您认为您可以详细说明一个非常简单的示例吗?不要担心,如果你不能,我会弄清楚的。
    • 正如文档所说“该项目可能被一个或多个订阅者丢弃”我如何重现这种行为?我只是第一次尝试 onNext 和 onSubscribe
    • @paul 我也对这样的例子感兴趣(现在试试)。如果我想出答案,会更新答案。
    • 谢谢,如果我成功了,我也会通知你
    • 可能尝试指定least value to the timeout to offer 可以在这里工作。
    【解决方案2】:

    SubmissionPublisher.offer 声明

    如果资源限制,该项目可能会被一个或多个订阅者丢弃 超出,在这种情况下,给定的处理程序(如果非空)是 调用,如果返回true,则重试一次。

    只是为了理解,在你的两个电话中

    publisher.offer("item", (subscriber, value) -> true); // the handler would be invoked
    
    publisher.offer("item", (subscriber, value) -> false); // the handler wouldn't be invoked
    

    publisher 仍会将给定项目发布给其当前的每个订阅者。这发生在您当前的场景中。


    正如文档建议的那样,在资源限制方面,验证您提供的处理程序是否通过尝试重现被调用的场景非常困难:

    如果资源限制,该项目可能会被一个或多个订阅者丢弃 超出,在这种情况下,给定的处理程序(如果非空)是 调用,如果返回true,重试一次。

    但是,您可以尝试使用将超时设置为基本最小值的项目删除 offer​(T item, long timeout, TimeUnit unit, BiPredicate&lt;Flow.Subscriber&lt;? super T&gt;,? super T&gt; onDrop)的重载方法

    timeout - 多长时间之前等待任何订阅者的资源 放弃,以单位为单位

    unit - 一个 TimeUnit 确定如何 解释超时参数

    由于offer 方法可能会丢弃项目(立即或使用有界超时),这将提供插入处理程序然后重试的机会。

    【讨论】:

    • 感谢您的信息 ;)
    猜你喜欢
    • 1970-01-01
    • 2023-03-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-05-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多