【问题标题】:Subscribing to Hot Observable source, then signalling on the same source订阅 Hot Observable 源,然后在同一源上发出信号
【发布时间】:2020-07-03 07:12:05
【问题描述】:

Tl;博士:

在我订阅了同一个 hot observable 之后,如何在 HotObservable 上发出 onNext 信号?

加长版:

我的目标是创建 REST 端点,它会发出事件,然后在返回响应之前等待(这本身就是一个事件)——基本上是带有命令/响应的事件驱动的 REST 端点。

由于我还没有找到任何解决方案,我决定通过反应式 Java 来完成,通过将 spring 事件映射到 DirectProcessor 发布者(Hot observable)。我有:

https://github.com/Venthe/Exploratory-Projects/tree/reactive/reactive

@Slf4j
@Component
class EventDriver {
    private static DirectProcessor<Object> EVENTS = DirectProcessor.create();

    @EventListener(Object.class)
    public void onEvent(Object event) {
        log.info(MessageFormat.format("Event {0} is ready to be pushed into processor.", event.toString()));
        EVENTS.onNext(event);
    }

    // Not important, but it does show that my events are captured & processed via 'normal' subscription
    @EventListener(ApplicationStartedEvent.class)
    public void logger() {
        EVENTS.map(Object::toString).subscribe(e -> log.info(MessageFormat.format("Event {0} received by processor.", e)));
    }

    public Flux<Object> getEvents() {
        log.info("Requesting processor.");
        return EVENTS
                .doOnEach(l -> log.info(MessageFormat.format("Processing subscribed event {0}", l.toString())));
    }
}
@RestController
@RequiredArgsConstructor
class ReservationRestController {
    private final ApplicationEventPublisher eventDispatcher;
    private final EventDriver eventDriver;

    @GetMapping(produces = MediaType.TEXT_EVENT_STREAM_VALUE, value = "/wait-for-event/{name}")
    public Flux<String> performAction(@PathVariable String name) {
        return eventDriver.getEvents()
                .doOnSubscribe(s -> eventDispatcher.publishEvent(new MyEvent(name)))
                .map(Object::toString)
                .log();
    }

    @AllArgsConstructor
    @NoArgsConstructor
    @Data
    class MyEvent {
        private String name;
    }
}

我的理解是:当我打开 REST 端点 \wait-for-event\any 时,响应式 Web 方法 performAction 应该首先订阅事件,然后调度事件(因此 - 触发管道),因为 doOnSubscribe(添加行为(副作用)在 {@link Flux} 完成订阅时触发)

不幸的是,事件卡住了 - 订阅没有看到我的事件,这是日志

 : Requesting processor.
 : Event ReservationRestController.MyEvent(name=as3) is ready to be pushed into processor.
 : Event ReservationRestController.MyEvent(name=as3) received by processor.
 : | onSubscribe([Fuseable] FluxOnAssembly.OnAssemblySubscriber)
 : | request(1)

据我了解,我是在实际调用 onSubscribe 之前调度的。

当然,之后发生的任何其他事件都会被正确捕获。

【问题讨论】:

    标签: java spring project-reactor


    【解决方案1】:

    一段时间后我又回到了这个问题;解决方案是先订阅事件,然后使用调度程序压缩:

    public String wtfperformActionPathVariable String name) {
            return Mono.from(
                    eventDriver.getEvents()
                            .map(Object::toString)
                            .filter(b -> b.contains(name))
                            .log("Event source")
            )
                    .zipWith(Mono.just(name)
                            .map(MyEvent::new)
                            .doOnNext(eventDispatcher::publishEvent)
                            .log("Publisher"))
                    .map(Tuple2::getT1)
                    .log("Result")
                    .block();
        }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-11-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多