从数据组装的演进,聊聊命令式与声明式的思维转变

# 技术# 设计
1.5k
4 min

“数据组装”(Data Enrichment)大概是每个写过业务 CRUD 的后端开发都绕不开的话题。随便查个订单列表 List<Order>,为了给前端展示,马上就得跟上回填“用户信息”、“商品信息”、“优惠券信息”。

其实大家都知道最直接的做法:“查出主表,然后批量查关联表,最后拼起来”。但在稍微有点规模的系统里,特别是当你开始考虑并发和响应延迟时,这种写代码的习惯会让人非常痛苦。

我最近在重构项目里的数据回填逻辑,顺便重新审视了下 DataLoaderEnrichmentService 的设计。这不仅仅是把代码抽成公共方法的问题,而是从“命令式(Imperative)”向“声明式(Declarative)”思维转变的过程。

为什么“命令式”写法让人崩溃?

大多数人一开始写数据组装,代码往往长这样:

List<Order> orders = orderRepository.findRecent();

// 1. 组装用户信息
Set<Long> userIds = orders.stream()
    .map(Order::getUserId)
    .filter(Objects::nonNull)
    .collect(Collectors.toSet());
if (!userIds.isEmpty()) {
    Map<Long, User> userMap = userLoader.load(userIds);
    orders.forEach(order -> {
        if (order.getUserId() != null) {
            order.setUser(userMap.get(order.getUserId()));
        }
    });
}

// 2. 组装商品信息
Set<Long> productIds = orders.stream()
    .map(Order::getProductId)
    .filter(Objects::nonNull)
    .collect(Collectors.toSet());
if (!productIds.isEmpty()) {
    Map<Long, Product> productMap = productLoader.load(productIds);
    orders.forEach(order -> {
        if (order.getProductId() != null) {
            order.setProduct(productMap.get(order.getProductId()));
        }
    });
}

这段代码能跑,逻辑也直白。但每次写这种代码,我都觉得自己在做无情的打字机:收集 Key、过滤 Null、批量查询、再写个双重循环回填。如果一个对象有 5 个关联实体,这段长得几乎一模一样的代码就得抄 5 遍。

真正让人崩溃的是加并发。为了缩短接口响应时间,我们往往得把“加载用户”和“加载商品”改成并行执行。于是,简单的逻辑被裹上了一层厚厚的并发原语:

// 为了并发,我们需要引入线程池和 CompletableFuture
CompletableFuture<Void> userFuture = CompletableFuture.runAsync(() -> {
    Set<Long> userIds = orders.stream().map(Order::getUserId).filter(Objects::nonNull).collect(Collectors.toSet());
    if (!userIds.isEmpty()) {
        Map<Long, User> userMap = userLoader.load(userIds);
        orders.forEach(order -> {
            if (order.getUserId() != null) order.setUser(userMap.get(order.getUserId()));
        });
    }
}, customExecutor);

CompletableFuture<Void> productFuture = CompletableFuture.runAsync(() -> {
    Set<Long> productIds = orders.stream().map(Order::getProductId).filter(Objects::nonNull).collect(Collectors.toSet());
    if (!productIds.isEmpty()) {
        Map<Long, Product> productMap = productLoader.load(productIds);
        orders.forEach(order -> {
            if (order.getProductId() != null) order.setProduct(productMap.get(order.getProductId()));
        });
    }
}, customExecutor);

// 等待所有任务完成
CompletableFuture.allOf(userFuture, productFuture).join();

这代码看起来就很“硬”。原本清晰的业务意图,全被 CompletableFuture、线程池配置和异常处理给埋了。只要哪个小伙伴稍微手抖,少写了个异常捕获或者用错了线程池,线上问题就跟着来了。

换个思路:只说“要什么”

其实我们每天都在写最纯粹的声明式代码——SQL。

想查已支付订单的用户名,在代码里我们得写一堆循环匹配,但在数据库里,一句 SQL 就搞定了:

SELECT o.id, u.username
FROM orders o
JOIN users u ON o.user_id = u.id
WHERE o.status = 'PAID';

写 SQL 的时候,我们只说“我要关联表里的这个字段”,至于数据库底层是怎么用 B+ 树做索引的,是用 Nested Loop 还是 Hash Join,我们根本不关心。这种脏活累活,引擎全包了。

在 Java 的数据组装里,我们也可以套用这个思路。

在项目里,我们先用 DataLoader 把“怎么取数据”这件小事给抽象出来:

@FunctionalInterface
public interface DataLoader<K, E> {
    Map<K, E> load(Set<K> keys);
}

然后,才是重头戏。我们用 EnrichmentService 来声明组装规则,而不是一步步去写循环:

enrichmentService.target(orders)
    .enrich(UserLoader.class, Order::getUserId, Order::setUser)
    .enrich(ProductLoader.class, Order::getProductId, Order::setProduct)
    .parallel() // 一行代码,开启并发
    .execute();

就是这么短。enrich 告诉框架用什么 Loader,怎么提取 Key,又怎么塞回 Target。至于底层怎么并发、怎么拼装参数,开发者连看都不用看。

扒开 EnrichmentService 看看

其实 EnrichmentService 没用什么黑魔法。阅读它的源码,你会发现它只是把那些又脏又累的控制流代码给藏起来了。

在内部,每次 .enrich(...) 都会注册成一个 EnrichmentTask

private record EnrichmentTask<T, K, E>(
        Class<? extends DataLoader<K, E>> loaderClass,
        Function<T, K> keyExtractor,
        BiConsumer<T, E> dataSetter) {

    LoadedData<T, K, E> loadData(Collection<T> targets, EnrichmentService service) {
        Set<K> keys = targets.stream()
                .map(keyExtractor)
                .filter(Objects::nonNull)
                .collect(Collectors.toSet());

        Map<K, E> dataMap = Map.of();
        if (!keys.isEmpty()) {
            DataLoader<K, E> loader = service.getLoader(loaderClass);
            dataMap = loader.load(keys);
        }
        return new LoadedData<>(this, dataMap);
    }
}

这种设计的爽点在于执行策略的隔离。如果用传统的命令式写法,从串行改并行,你需要重写一多半的代码。但由于我们已经把“要干嘛”和“怎么干”剥离了,EnrichmentExecutor 内部可以随意切换底层引擎。

比如它的并行实现:

private void executeInParallel() {
    final Executor activeExecutor = this.overrideExecutor != null ? this.overrideExecutor
            : this.service.defaultExecutor;

    var futures = tasks.stream()
            .map(task -> CompletableFuture.supplyAsync(() -> task.loadData(targets, service), activeExecutor))
            .toList();

    CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();

    futures.stream()
            .map(CompletableFuture::join)
            .forEach(loadedData -> loadedData.applySetter(targets));
}

有了 Java 21+ 之后,甚至可以直接在链式调用里声明 .withVirtualThreads()。业务代码一行不改,底层就无缝切到了虚拟线程。不用再去写一堆协程调度的脏代码,这感觉很奇妙。

看看别人是怎么做的

在 Java 圈子里解决“N+1”和数据装配问题,早有几个非常成熟的轮子,我也拿来和我们的设计做过横向对比。

GraphQL-Java / java-dataloader Facebook DataLoader 的 Java 官方移植版。这东西在 GraphQL 社区用得很广。它强制你写 BatchLoader<K, V>,思想跟我们极其一致。不过,它主要靠请求级缓存和事件回调驱动,跟 GraphQL 的执行引擎绑得很死。如果在普通的 Spring Boot 接口或者 MQ 消费者里硬拉它进来,开发体验实在有点繁琐。

Crane4j 国内开源的框架。它重度依赖注解(比如 @Assemble),还有 Spring AOP 和反射。功能包罗万象,连全局缓存和数据转换都给你做好了。是个大而全的重型武器。相比之下,我们的 EnrichmentService 就是个轻量的纯 Java 工具。我个人偏好编译期能看到明确的类型调用,而不是靠满天飞的反射和注解去猜运行时的行为。

命令式 vs 声明式

维度 命令式 (Imperative) 声明式 (Declarative)
你在想什么 控制流、中间变量怎么存、并发怎么合并 组装规则是什么、字段的映射关系
可读性 随着字段增多,业务逻辑变成一坨面条 一行代码配一个规则,一眼望穿
重构成本 牵一发而动全身(尤其是改并发) 一键切换执行引擎(串行/并行/协程)

结语

写业务代码写久了,容易陷入一种思维惯性:不管什么需求,上来就先写一个 for 循环。

但其实很多时候,代码写得乱,并不是需求有多复杂,而是我们把“业务规则”和“控制流程”紧紧揉在了一起。框架和架构设计的意义,就在于把那些容易写错的、机械性的步骤沉入水底。

声明式编程不是万能的,但当你需要处理大量雷同的业务逻辑时,换个角度,只告诉系统“你想要什么”,代码真的会干净很多。

附录:完整源码

DataLoader.java

import java.util.Map;
import java.util.Set;

@FunctionalInterface
public interface DataLoader<K, E> {

    // 根据一批Key,从数据源批量获取数据实体。
    Map<K, E> load(Set<K> keys);
}

EnrichmentService.java

import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.concurrent.Executors;
import java.util.function.BiConsumer;
import java.util.function.Function;
import java.util.stream.Collectors;

public class EnrichmentService {

    private final Map<Class<? extends DataLoader>, DataLoader<?, ?>> loaderMap;
    private final Executor defaultExecutor;

    // 构造函数
    public EnrichmentService(List<DataLoader<?, ?>> loaders, Executor defaultExecutor) {
        this.loaderMap = loaders.stream()
                .collect(Collectors.toMap(DataLoader::getClass, l -> l, (l1, l2) -> l1));
        this.defaultExecutor = Objects.requireNonNull(defaultExecutor, "默认的执行器不能为空");
    }

    // 创建包含平台线程池的服务实例
    public static EnrichmentService createStandalone(List<DataLoader<?, ?>> loaders) {
        return new EnrichmentService(loaders, Executors.newCachedThreadPool());
    }

    // 创建包含虚拟线程池的服务实例
    public static EnrichmentService createStandaloneWithVirtualThreads(List<DataLoader<?, ?>> loaders) {
        return new EnrichmentService(loaders, Executors.newVirtualThreadPerTaskExecutor());
    }

    // 创建执行器
    public <T> EnrichmentExecutor<T> target(Collection<T> targets) {
        return new EnrichmentExecutor<>(targets, this);
    }

    @SuppressWarnings("unchecked")
    protected <K, E> DataLoader<K, E> getLoader(Class<? extends DataLoader<K, E>> loaderClass) {
        DataLoader<K, E> loader = (DataLoader<K, E>) loaderMap.get(loaderClass);
        if (loader == null) {
            throw new IllegalArgumentException("未找到类型为 " + loaderClass.getName() + " 的DataLoader实例或Bean");
        }
        return loader;
    }

    // 流式执行器,支持多种模式
    public static class EnrichmentExecutor<T> {
        private final Collection<T> targets;
        private final EnrichmentService service;
        private final List<EnrichmentTask<T, ?, ?>> tasks = new ArrayList<>();
        private Executor overrideExecutor = null;
        private ExecutionMode mode = ExecutionMode.SEQUENTIAL;

        private enum ExecutionMode {
            SEQUENTIAL, PARALLEL
        }

        public EnrichmentExecutor(Collection<T> targets, EnrichmentService service) {
            this.targets = targets;
            this.service = service;
        }

        // 注册任务
        public <K, E> EnrichmentExecutor<T> enrich(Class<? extends DataLoader<K, E>> loaderClass,
                Function<T, K> keyExtractor, BiConsumer<T, E> dataSetter) {
            this.tasks.add(new EnrichmentTask<>(loaderClass, keyExtractor, dataSetter));
            return this;
        }

        // 串行
        public EnrichmentExecutor<T> sequentially() {
            this.mode = ExecutionMode.SEQUENTIAL;
            return this;
        }

        // 并行
        public EnrichmentExecutor<T> parallel() {
            this.mode = ExecutionMode.PARALLEL;
            return this;
        }

        // 使用虚拟线程
        public EnrichmentExecutor<T> withVirtualThreads() {
            this.mode = ExecutionMode.PARALLEL;
            this.overrideExecutor = Executors.newVirtualThreadPerTaskExecutor();
            return this;
        }

        // 使用自定义线程池
        public EnrichmentExecutor<T> usingExecutor(Executor executor) {
            this.mode = ExecutionMode.PARALLEL;
            this.overrideExecutor = Objects.requireNonNull(executor, "自定义执行器不能为空");
            return this;
        }

        // 执行任务
        public void execute() {
            if (targets == null || targets.isEmpty() || tasks.isEmpty()) {
                return;
            }

            switch (this.mode) {
                case SEQUENTIAL -> executeSequentially();
                case PARALLEL -> executeInParallel();
            }
        }

        private void executeSequentially() {
            for (var task : tasks) {
                var loadedData = task.loadData(targets, service);
                loadedData.applySetter(targets);
            }
        }

        private void executeInParallel() {
            final Executor activeExecutor = this.overrideExecutor != null ? this.overrideExecutor
                    : this.service.defaultExecutor;

            var futures = tasks.stream()
                    .map(task -> CompletableFuture.supplyAsync(() -> task.loadData(targets, service), activeExecutor))
                    .toList();

            CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();

            futures.stream()
                    .map(CompletableFuture::join)
                    .forEach(loadedData -> loadedData.applySetter(targets));
        }
    }

    // 内部辅助记录类
    private record EnrichmentTask<T, K, E>(
            Class<? extends DataLoader<K, E>> loaderClass,
            Function<T, K> keyExtractor,
            BiConsumer<T, E> dataSetter) {
        LoadedData<T, K, E> loadData(Collection<T> targets, EnrichmentService service) {
            Set<K> keys = targets.stream()
                    .map(keyExtractor)
                    .filter(Objects::nonNull)
                    .collect(Collectors.toSet());

            Map<K, E> dataMap = Map.of();
            if (!keys.isEmpty()) {
                DataLoader<K, E> loader = service.getLoader(loaderClass);
                dataMap = loader.load(keys);
            }
            return new LoadedData<>(this, dataMap);
        }
    }

    private record LoadedData<T, K, E>(EnrichmentTask<T, K, E> task, Map<K, E> dataMap) {
        void applySetter(Collection<T> targets) {
            if (dataMap == null || dataMap.isEmpty())
                return;

            for (T target : targets) {
                K key = task.keyExtractor().apply(target);
                if (key != null) {
                    E entity = dataMap.get(key);
                    if (entity != null) {
                        task.dataSetter().accept(target, entity);
                    }
                }
            }
        }
    }
}