You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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()));
}

注意事项

  1. 排序要求:必须确保两个流严格按共享键升序排列(如果是降序,需修改比较逻辑),否则匹配会出错
  2. 性能优势:手动归并的方式是O(n+m)时间复杂度,不需要缓存全量数据,适合大数据量场景
  3. 小数据量场景:如果数据量不大,也可以先将两个流收集为Map<ClientId, List<...>>,再遍历匹配,但效率不如归并方式

内容的提问来源于stack exchange,提问作者Adrian I

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 04:22:52