如何在Reactor中基于共享键连接不同对象的Publisher(类SQL连接)
基于Reactor实现已排序Publisher的SQL风格JOIN操作
针对两个已按共享键排序的Publisher(以Reactor的Flux为例),我们可以用归并匹配的方式实现类似SQL的内连接、左连接、右连接,这种方式不需要缓存全量数据,效率更高。以下是具体实现(以Client和Invoice为例,二者共享clientId键):
先定义实体类
// 客户端实体 record Client(Long clientId, String name) {} // 发票实体 record Invoice(Long clientId, BigDecimal amount) {} // 连接结果实体 record ClientInvoice(Long clientId, Client client, Invoice invoice) {}
内连接(Inner Join)
只保留两边都存在匹配键的记录,逻辑是按归并排序的方式遍历两个流,匹配相同clientId的记录:
public Flux<ClientInvoice> innerJoin(Flux<Client> clients, Flux<Invoice> invoices) { return Flux.create(sink -> { AtomicReference<Client> currentClient = new AtomicReference<>(); AtomicReference<Invoice> currentInvoice = new AtomicReference<>(); Subscription clientSub = null; Subscription invoiceSub = null; // 订阅客户端流 clientSub = clients.subscribe( client -> { currentClient.set(client); // 尝试匹配当前发票 while (currentInvoice.get() != null) { int cmp = client.clientId().compareTo(currentInvoice.get().clientId()); if (cmp == 0) { sink.next(new ClientInvoice(client.clientId(), client, currentInvoice.get())); currentInvoice.set(null); invoiceSub.request(1); } else if (cmp < 0) { // 当前客户端ID更小,等待发票流跟进 break; } else { // 当前发票ID更小,跳过该发票 currentInvoice.set(null); invoiceSub.request(1); } } }, sink::error, () -> sink.complete() ); // 订阅发票流 invoiceSub = invoices.subscribe( invoice -> { currentInvoice.set(invoice); // 尝试匹配当前客户端 while (currentClient.get() != null) { int cmp = currentClient.get().clientId().compareTo(invoice.clientId()); if (cmp == 0) { sink.next(new ClientInvoice(invoice.clientId(), currentClient.get(), invoice)); currentClient.set(null); clientSub.request(1); } else if (cmp > 0) { // 当前发票ID更小,等待客户端流跟进 break; } else { // 当前客户端ID更小,跳过该客户端 currentClient.set(null); clientSub.request(1); } } }, sink::error, () -> sink.complete() ); // 初始请求数据 clientSub.request(1); invoiceSub.request(1); }); }
左连接(Left Join)
保留所有客户端记录,无匹配发票的记录中invoice字段为null:
public Flux<ClientInvoice> leftJoin(Flux<Client> clients, Flux<Invoice> invoices) { return Flux.create(sink -> { AtomicReference<Client> currentClient = new AtomicReference<>(); AtomicReference<Invoice> currentInvoice = new AtomicReference<>(); Subscription clientSub = null; Subscription invoiceSub = null; boolean invoicesCompleted = false; clientSub = clients.subscribe( client -> { currentClient.set(client); // 跳过所有ID小于当前客户端的发票(无匹配) while (currentInvoice.get() != null && client.clientId() > currentInvoice.get().clientId()) { currentInvoice.set(null); invoiceSub.request(1); } // 匹配所有同ID的发票 boolean hasMatch = false; while (currentInvoice.get() != null && client.clientId().equals(currentInvoice.get().clientId())) { sink.next(new ClientInvoice(client.clientId(), client, currentInvoice.get())); hasMatch = true; currentInvoice.set(null); invoiceSub.request(1); } // 如果发票流已结束且无匹配,输出客户端+null if (!hasMatch && invoicesCompleted) { sink.next(new ClientInvoice(client.clientId(), client, null)); currentClient.set(null); clientSub.request(1); } // 如果发票ID大于当前客户端ID,说明无匹配,直接输出 else if (!hasMatch && currentInvoice.get() != null && currentInvoice.get().clientId() > client.clientId()) { sink.next(new ClientInvoice(client.clientId(), client, null)); currentClient.set(null); clientSub.request(1); } }, sink::error, () -> sink.complete() ); invoiceSub = invoices.subscribe( invoice -> { currentInvoice.set(invoice); while (currentClient.get() != null) { int cmp = currentClient.get().clientId().compareTo(invoice.clientId()); if (cmp == 0) { sink.next(new ClientInvoice(invoice.clientId(), currentClient.get(), invoice)); currentClient.set(null); clientSub.request(1); currentInvoice.set(null); invoiceSub.request(1); } else if (cmp < 0) { // 当前客户端ID更小,无匹配发票,输出客户端+null sink.next(new ClientInvoice(currentClient.get().clientId(), currentClient.get(), null)); currentClient.set(null); clientSub.request(1); } else { // 当前发票ID更小,跳过 currentInvoice.set(null); invoiceSub.request(1); break; } } }, sink::error, () -> { invoicesCompleted = true; // 发票流结束后,输出剩余未匹配的客户端 while (currentClient.get() != null) { sink.next(new ClientInvoice(currentClient.get().clientId(), currentClient.get(), null)); currentClient.set(null); clientSub.request(1); } sink.complete(); } ); clientSub.request(1); invoiceSub.request(1); }); }
右连接(Right Join)
保留所有发票记录,无匹配客户端的记录中client字段为null,可以直接复用左连接逻辑反转流即可:
public Flux<ClientInvoice> rightJoin(Flux<Client> clients, Flux<Invoice> invoices) { return leftJoin(invoices, clients) .map(result -> new ClientInvoice(result.clientId(), result.invoice(), result.client())); }
注意事项
- 排序要求:必须确保两个流严格按共享键升序排列(如果是降序,需修改比较逻辑),否则匹配会出错
- 性能优势:手动归并的方式是O(n+m)时间复杂度,不需要缓存全量数据,适合大数据量场景
- 小数据量场景:如果数据量不大,也可以先将两个流收集为
Map<ClientId, List<...>>,再遍历匹配,但效率不如归并方式
内容的提问来源于stack exchange,提问作者Adrian I
相关产品推荐
相关产品推荐

