Vik*_*ahl 5 java event-driven reactor project-reactor
我正在学习Reactor,并且想知道如何实现某种行为。假设我有一堆传入消息。每个消息都与某个实体相关联,并包含一些数据。
interface Message {
String getEntityId();
Data getData();
}
Run Code Online (Sandbox Code Playgroud)
与不同实体有关的消息可以并行处理。但是,与任何单个实体有关的消息必须一次处理一次,即,对于实体的消息2的处理要"abc"等到针对实体的消息1的处理"abc"完成后才能开始。在处理消息的过程中,应缓存该整段消息。其他实体的消息可以不受阻碍地进行。可以认为它是每个实体上运行这样的代码的线程:
public void run() {
for (;;) {
// Blocks until there's a message available
Message msg = messageQueue.nextMessageFor(this.entityId);
// Blocks until processing is finished
processMessage(msg);
}
}
Run Code Online (Sandbox Code Playgroud)
我该如何使用React做到这一点而又不会受到阻碍?总消息速率可能很高,但是每个实体的消息速率会非常低。实体集可能非常大,因此不一定事先知道。
我猜可能看起来像这样,但我不知道。
{
incomingMessages()
.groupBy(Message::getEntityId)
.flatMap(entityStream -> entityStream
/* ... */
.map(msg -> /* process the message */)))
/* ... */
}
public static Stream<Message> incomingMessages() { /* ... */ }
Run Code Online (Sandbox Code Playgroud)
使用ProjectReactor可以这样解决:
@Test
public void testMessages() {
Flux.fromStream(incomingMessages())
.groupBy(Message::getEntityId)
.map(g -> g.publishOn(Schedulers.newParallel("groupByPool", 16))) //create new publisher for groups of messages
.subscribe( //create consumer for main stream
stream ->
stream.subscribe(this::processMessage) // create consumer for group stream
);
}
public Stream<Message> incomingMessages() {
return IntStream.range(0, 100).mapToObj(i -> new Message(i, i % 10));
}
public void processMessage(Message message) {
System.out.println(String.format("Message: %s processed by the thread: %s", message, Thread.currentThread().getName()));
}
private static class Message {
private final int id;
private final int entityId;
public Message(int id, int entityId) {
this.id = id;
this.entityId = entityId;
}
public int getId() {
return id;
}
public int getEntityId() {
return entityId;
}
@Override
public String toString() {
return "Message{" +
"id=" + id +
", entityId=" + entityId +
'}';
}
}
Run Code Online (Sandbox Code Playgroud)
我认为类似的解决方案可以在RxJava中
我们的项目中也遇到了同样的问题。具有相同 id 的实体必须按顺序处理,但具有不同 id 的实体可以并行处理。
解决方案非常简单。我们开始使用 concatMap,而不是使用 flatMap。来自 concatMap 的文档:
* Transform the elements emitted by this {@link Flux} asynchronously into Publishers,
* then flatten these inner publishers into a single {@link Flux}, sequentially and
* preserving order using concatenation.
Run Code Online (Sandbox Code Playgroud)
代码示例:
public void receive(Flux<Data> data) {
data
.groupBy(Data::getPointID)
.flatMap(service::process)
.onErrorContinue(Logging::logError)
.subscribe();
}
Run Code Online (Sandbox Code Playgroud)
工艺方法:
Flux<SomeEntity> process(Flux<Data> dataFlux) {
return dataFlux
.doOnNext(Logging::logReceived)
.concatMap(this::proceedDefinitionsSearch)
.doOnNext(Logging::logDefSearch)
.flatMap(this::processData)
.doOnNext(Logging::logDataProcessed)
.concatMap(repository::save)
.doOnNext(Logging::logSavedEntity);
}
Run Code Online (Sandbox Code Playgroud)