从数据组装的演进,聊聊命令式与声明式的思维转变
“数据组装”(Data Enrichment)大概是每个写过业务 CRUD 的后端开发都绕不开的话题。随便查个订单列表 List<Order>,为了给前端展示,马上就得跟上回填“用户信息”、“商品信息”、“优惠券信息”。
其实大家都知道最直接的做法:“查出主表,然后批量查关联表,最后拼起来”。但在稍微有点规模的系统里,特别是当你开始考虑并发和响应延迟时,这种写代码的习惯会让人非常痛苦。
我最近在重构项目里的数据回填逻辑,顺便重新审视了下 DataLoader 和 EnrichmentService 的设计。这不仅仅是把代码抽成公共方法的问题,而是从“命令式(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 循环。
但其实很多时候,代码写得乱,并不是需求有多复杂,而是我们把“业务规则”和“控制流程”紧紧揉在了一起。框架和架构设计的意义,就在于把那些容易写错的、机械性的步骤沉入水底。
声明式编程不是万能的,但当你需要处理大量雷同的业务逻辑时,换个角度,只告诉系统“你想要什么”,代码真的会干净很多。
附录:完整源码
import java.util.Map;
import java.util.Set;
@FunctionalInterface
public interface DataLoader<K, E> {
// 根据一批Key,从数据源批量获取数据实体。
Map<K, E> load(Set<K> keys);
}
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);
}
}
}
}
}
}