【发布时间】: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