【问题标题】:I am losing messages using spring-reactor, what is wrong with my setup?我正在使用 spring-reactor 丢失消息,我的设置有什么问题?
【发布时间】:2013-12-05 12:52:34
【问题描述】:

我想我应该研究一下 Pivotal 新发布的反应器框架,用于我正在编写的一个简单程序,该程序需要一些多线程来及时完成。

我编写了以下示例项目来了解框架并使用它来了解它的使用方式:

Main.java:

package reactortest;

import org.springframework.context.annotation.AnnotationConfigApplicationContext;

public class Main { 
    public static void main(String[] args) throws InterruptedException {
        try(AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(MainConfiguration.class)) {
            MyProducer producer = context.getBean(MyProducer.class);
            producer.run();
        }
    }
}

MyProducer.java:

package reactortest;

import java.util.concurrent.CountDownLatch;

import reactor.core.Reactor;
import reactor.event.Event;

public class MyProducer {
    private final Reactor reactor;
    private final Integer messagesToPrint;
    private final CountDownLatch countDownLatch;

    public MyProducer(final Reactor reactor, final Integer messagesToPrint, CountDownLatch countDownLatch) {
        this.reactor = reactor;
        this.messagesToPrint = messagesToPrint;
        this.countDownLatch = countDownLatch;
    }

    public void run() throws InterruptedException {
        for(int i = 0; i < messagesToPrint; ++i) {
            reactor.notify(Event.wrap("String event: " + i));
        }

        countDownLatch.await();
        System.out.println("Finished. Remaining count is: " + countDownLatch.getCount());
    }
}

MyConsumer.java:

package reactortest;

import java.util.concurrent.CountDownLatch;

import reactor.event.Event;
import reactor.function.Consumer;

public class MyConsumer implements Consumer<Event<String>> {
    private final CountDownLatch countDownLatch;

    public MyConsumer(CountDownLatch countDownLatch) {
        this.countDownLatch = countDownLatch;
    }

    @Override
    public void accept(Event<String> message) {
        System.out.println(message);
        countDownLatch.countDown();
    }
}

最后, MainConfiguration.java:

package reactortest;

import java.util.concurrent.CountDownLatch;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import reactor.core.Environment;
import reactor.core.Reactor;
import reactor.core.spec.Reactors;
import reactor.spring.context.config.EnableReactor;

@Configuration
@EnableReactor
public class MainConfiguration {
    private final Integer MESSAGESTOPRINT = 10;

    @Autowired private Environment environment;

    @Bean
    public CountDownLatch countDownLatch() {
        CountDownLatch countDownLatch = new CountDownLatch(MESSAGESTOPRINT);
        return countDownLatch;
    }

    @Bean
    public Reactor reactor() {
        Reactor reactor = Reactors.reactor().env(environment).dispatcher(Environment.THREAD_POOL).randomEventRouting().get();
        reactor.on(consumer());
        return reactor;
    }

    @Bean
    public MyProducer producer() {
        MyProducer producer = new MyProducer(reactor(), MESSAGESTOPRINT, countDownLatch());
        return producer;
    }

    @Bean
    public MyConsumer consumer() {
        MyConsumer consumer = new MyConsumer(countDownLatch());
        return consumer;
    }
}

我的问题是程序永远不会结束。消费者每次运行时也会打印出不同的信息。从连续三个运行它打印:

1st run:
Event{id=null, headers=null, replyTo=null, data=String event: 0}
Event{id=null, headers=null, replyTo=null, data=String event: 1}
Event{id=null, headers=null, replyTo=null, data=String event: 7}
Event{id=null, headers=null, replyTo=null, data=String event: 8}

2nd run:
Event{id=null, headers=null, replyTo=null, data=String event: 0}
Event{id=null, headers=null, replyTo=null, data=String event: 1}
Event{id=null, headers=null, replyTo=null, data=String event: 5}
Event{id=null, headers=null, replyTo=null, data=String event: 6}
Event{id=null, headers=null, replyTo=null, data=String event: 9}

3rd run:
Event{id=null, headers=null, replyTo=null, data=String event: 2}
Event{id=null, headers=null, replyTo=null, data=String event: 4}
Event{id=null, headers=null, replyTo=null, data=String event: 6}

似乎我一定错过了一些非常明显的东西,因为除了这个是 javaconfig 而不是配置的注释,并且没有与外界进行任何交互之外,我看不出这与示例 here 有何不同。

【问题讨论】:

    标签: java spring reactor


    【解决方案1】:

    在问这个问题时,我正在改进代码,它最终奏效了(那里有一些 great rubber ducking)。我认为与其删除我的问题,不如将其发布,以防其他人遇到同样的问题。

    上面代码的问题是在设置reactor时调用randomEventRouting(),当设置这个标志时它随机选择消费者路由。因为我没有设置特定的选择器/键来定义要调度的消费者,并且由于在没有提供密钥时所有消费者都匹配,所以我假设在幕后设置了一个默认消费者,它正在传递我的一些事件.

    更改 reactor.on() 以接受选择器:

    reactor.on(Selectors.$(selector()), consumer());
    

    选择器在哪里:

    @Bean
    public String selector() {
        String selector = "My very special event";
        return selector;
    }
    

    并将此密钥注入生产者,并在调用 reactor.notify() 时使用它:

    reactor.notify(selector, Event.wrap("String event: " + i));
    

    按预期工作。

    我想这是一个非常极端的情况,因为大多数用户会(并且应该)定义键,但你永远不知道。 :)

    【讨论】:

      猜你喜欢
      • 2015-03-28
      • 1970-01-01
      • 1970-01-01
      • 2023-03-14
      • 1970-01-01
      • 2016-06-24
      • 1970-01-01
      • 2012-11-06
      • 1970-01-01
      相关资源
      最近更新 更多