There is no fork:Java Fetch 如何从声明式查询中发现并发
Fetch 把业务查询写成可组合的程序描述,再自动发现独立请求、合并批量查询,并沿着数据依赖分轮执行。
假设要根据一个年级 ID,查出这个年级的所有班级和所有学生。一种自然的写法是复用已有的小查询:
Fetch<List<GradeClassStudent>> program =
new GradeById(gradeId).toFetch()
.flatMap(grade ->
concatM(grade.getClasses(), classId ->
new ClassById(classId).toFetch()
.flatMap(clazz ->
mapM(clazz.getStudents(), studentId ->
new StudentById(studentId).toFetch()
.map(student ->
new GradeClassStudent(
grade, clazz, student))))));这段代码的结构和业务问题几乎一样:先找年级,再找年级里的班级,最后找每个班级里的学生。每一个函数只负责一小块查询,可以单独理解,也可以继续组合成更大的查询。
但如果按照普通 Java 代码的字面顺序执行,问题马上就会出现:
1. 查询一次年级;
2. 对每一个班级分别查询一次班级;
3. 对每一个学生分别查询一次学生。
假设一个年级有 3 个班,每个班有 30 个学生,就可能产生 1 + 3 + 90 次查询。这是典型的 N+1 问题。
传统解决方式是手工破坏这段自然结构:先遍历年级收集所有 classId,调用一次 classByIds;再从所有班级收集 studentId,调用一次 studentByIds;最后建立几个 Map,把结果重新组装回原来的树。
查询次数变少了,但代价是业务代码开始负责数据库的执行计划。原来可以复用的小查询被拆开,调用者必须知道哪些 ID 应该提前收集、哪些查询可以合并、哪些结果要转成 Map。类似的问题在不同页面重复出现,开发者不断手写一次性的查询优化器。
Fetch 想消除的就是这个取舍:业务代码保持声明式和模块化,执行器把它自动转换成批量、可并发的执行计划。
这也是 fetch intro 和 fetch in erp 里记录的实际出发点。
本文讨论的 Java 实现是 schneiderlin/fetch,它移植了论文 There is no fork: an abstraction for efficient, concurrent, and concise data access 中最核心的计算模型。
声明式查询留下了优化空间
普通命令式代码同时规定了两件事:
- 要得到什么结果;
- 每一步必须以什么顺序执行。
例如:
User user = userRepository.getById(userId);
List<Order> orders = orderRepository.listByUserId(userId);
Dashboard dashboard = new Dashboard(user, orders);即使两个查询互不依赖,解释器也必须先执行 getById,再执行 listByUserId。顺序已经被程序员写死。
Fetch 只描述结果之间的关系:
Fetch<Dashboard> dashboard = appCombine(
Dashboard::new,
new UserById(userId).toFetch(),
new OrdersByUserId(userId).toFetch());UserById 和 OrdersByUserId 都是完整、独立的 Fetch。appCombine 表示:拿到两个结果以后,用 Dashboard::new 把它们组合起来。至于先查哪一个、是否并发、怎样 batch,不属于业务语义。
再看列表查询:
Fetch<List<User>> users =
mapM(userIds, id -> new UserById(id).toFetch());字面上看,它像是对每一个 ID 发起一次查询。但 mapM 不会立即执行这些请求,而是先把它们组合成一个 Fetch<List<User>>。执行器看到同一轮中有很多 UserById,便可以把它们变成一次:
UserById.batchQuery(context, userIds);声明式程序规定了结果,保留了实现自由。执行器掌握整轮请求以后,才能做局部代码看不到的优化。这个思路与 SQL 类似:SQL 描述想要什么,查询优化器决定怎样得到它。参见 what fetch is trying to solve。
There is no fork
并发 Java 代码通常会显式出现这些东西:
executor.submit(...);
CompletableFuture.supplyAsync(...);
CompletableFuture.allOf(...);
future.get(...);业务代码不只要知道数据依赖,还要创建任务、选择线程池、保存 Future、等待结果和处理超时。更麻烦的是,只有当前函数里的任务容易被放进同一个 allOf;封装在其他模块里的查询很难参与统一调度。
Fetch 程序里没有 fork。不是因为没有并发,而是因为程序员不再显式发起和编排并发。
Fetch<A> fetch = request.toFetch();这里没有立即访问数据库。Fetch<A> 也不是已经提交给线程池的 Future<A>。它描述的是一段最终会得到 A、但在取得外部数据时可以暂停、之后可以恢复的程序。
可以把核心类型近似理解成:
Fetch<A> ≈ IO<Result<A>>
Result<A> = Done<A>
| Blocked<A>其中:
Done(value)表示当前计算已经完成;而:
Blocked(requests, continuation)表示当前计算缺少一组外部数据。requests 是现在已经发现的请求,continuation 是拿到数据后剩余的程序。
Continuation 就是“接下来还要做什么”。它使执行器可以在程序等待数据库时把剩余计算保存下来,先去探索其他分支,等数据回来再恢复。更基础的 Java 解释见 continuation explain in java。
一个请求怎样变成 Blocked
Java Fetch 的 dataFetch 大致做四件事:
static <K, A> Fetch<A> dataFetch(Request<K, A> request) {
IORef<Object> box = createBox();
BlockedRequest<K, A> blocked =
new BlockedRequest<>(box, request);
Fetch<A> continuation =
read(box).map(value -> new Done<>((A) value));
return new Fetch<>(
IO.value(new Blocked<>(List.of(blocked), continuation)));
}这里的 box 是请求和 continuation 之间的结果槽:
1. dataFetch 先返回 Blocked,把 request 和空 box 交给执行器;
2. resolver 执行查询,把结果写进 box;
3. 执行器恢复 continuation;
4. continuation 从 box 读取结果,变成 Done<A>。
所以 Fetch 并不是先建立一个完整的静态 AST,再单独运行优化器。它是一段 resumable computation:每次只向前解释到当前能够看见的 Done 或 Blocked,执行完请求后再继续。
执行器按 round 运行
runFetch 的核心循环非常短:
if (result instanceof Done<A> done) {
return IO.value(done.value());
}
if (result instanceof Blocked<A> blocked) {
return resolver.apply(blocked.requests())
.flatMap(ignored ->
runFetch(resolver, blocked.continuation()));
}一次完整执行不断重复三个阶段:
explore program
→ collect blocked requests
→ resolve and write results
→ resume continuation每次收集并执行一批请求叫作一个 round。Round 的数量主要由最长的数据依赖链决定,而不是由实体数量决定。
前面的年级例子会形成三轮:
Round 1 GradeById(grade1)
Round 2 ClassById(class1)
ClassById(class2) → classByIds([class1, class2, class3])
ClassById(class3)
Round 3 StudentById(student1)
StudentById(student2)
... → studentByIds([...])
StudentById(student90)第二轮必须等第一轮,因为只有拿到 Grade 才知道有哪些 classId。第三轮也必须等第二轮,因为只有拿到 Class 才知道有哪些 studentId。但每一轮内部的 fan-out 都能自动合并,不会因为有 90 个学生就执行 90 轮。
Applicative:同时探索互不依赖的分支
Fetch 怎样知道两个分支可以一起探索?答案在 Applicative 组合。
Java 版本中的 app 对应 Haskell 的 <*>:
app(
Fetch<Function<A, B>> function,
Fetch<A> argument)两个参数都是已经存在的 Fetch。右边的程序不需要左边产生的值才能被构造出来,因此 app 可以分别向前执行两边,观察它们当前是 Done 还是 Blocked:
Done(f) <*> Done(a) = Done(f(a))
Done(f) <*> Blocked(b) = Blocked(map(b.cont, f))
Blocked(a) <*> Done(b) = Blocked(app(a.cont, Done(b)))
Blocked(a) <*> Blocked(b) =
Blocked(
concat(a.requests, b.requests),
app(a.cont, b.cont));最后一种情况最重要。左右分支都遇到 Blocked 时,执行器不会先解决左边,而是把两边的 request 合并进同一个 Blocked,把两边的 continuation 也继续组合起来。
例如:
Blocked(Done(+1)) <*> Blocked(Done(1))同时探索两边后得到:
Blocked(Done(+1) <*> Done(1))
= Blocked(Done(2))只有一层 Blocked,所以只需要一个 round。
如果使用从 Monad 顺序推导出来的普通 ap,它只先观察左边:
Blocked(Done(+1)) <*> Blocked(Done(1))
= Blocked(Done(+1) <*> Blocked(Done(1)))
= Blocked(Blocked(Done(2)))两个本来独立的请求因此被放进两轮。它们并没有真实的数据依赖,只是解释器过早规定了观察顺序。
这里的 explore 不是执行完远程查询,而是尽可能向前计算,直到知道每个分支最外层是 Done 还是 Blocked。Applicative 的结构证明了两个 Fetch 可以独立探索,执行器才能把它们放进同一轮。
这也解释了为什么 applicative 不只是抽象数学工具。在 Fetch 里,Applicative 直接携带了调度信息。
Monad:保留真正的数据依赖
flatMap 对应 Monad 的 bind:
Fetch<B> flatMap(Function<A, Fetch<B>> next)这里的右边不是一个现成的 Fetch<B>,而是:
A -> Fetch<B>例如:
new GradeById(gradeId).toFetch()
.flatMap(grade ->
mapM(grade.getClasses(),
classId -> new ClassById(classId).toFetch()));在 Grade 返回以前,我们不知道 grade.getClasses() 是什么,也就无法构造右边的 ClassById 请求。执行器只能先得到 Grade,再调用 continuation 产生下一轮计算。
因此可以把两者理解成:
Applicative / appCombine
两边的 Fetch 已经存在
→ 表达独立关系
→ 可以同时 explore
Monad / flatMap
右边要用左边的值才能构造
→ 表达数据依赖
→ 必须分 round如果两个查询实际上独立,却写成 flatMap:
fetchA.flatMap(a ->
fetchB.map(b -> combine(a, b)));Fetch 不会分析 lambda 的自由变量,再猜测 fetchB 是否真的依赖 a。程序员需要用 appCombine 明确表达独立性:
appCombine(this::combine, fetchA, fetchB);这不是为了迎合函数式编程风格,而是在告诉执行器哪里存在可以安全合并的控制流。
mapM:把 N 个小查询变成一个 batch
mapM 是列表查询能够自动 batch 的关键 helper:
mapM(ids, id -> new UserById(id).toFetch())它先把每个 ID 映射成 Fetch<User>,再用 Applicative 组合这些 Fetch,而不是用 flatMap 串行 fold。因此所有 UserById 都能在同一轮被探索出来。
Resolver 收到这一轮的 BlockedRequest 后:
1. 按具体 Request class 分组;
2. 收集每组的 ID;
3. 对 ID 去重;
4. 反射调用 Request class 上的静态 batchQuery;
5. 把 batch 返回的 Map 分发到每个请求的 result box;
6. 恢复组合后的 continuation。
例如 Request 可以这样定义:
public final class UserById implements Request<Long, User> {
private final long id;
@Override
public Long getId() {
return id;
}
public static Map<Long, User> batchQuery(
FetchContext context,
List<Long> ids) {
UserRepository repository =
context.getBean("users", UserRepository.class);
return repository.findByIds(ids)
.toMap(User::id, Function.identity());
}
}Request 描述单个、容易组合的查询语义,batchQuery 描述这一类请求的批量物理执行方式。业务代码面向前者编程,resolver 在运行时选择后者。
Blog 程序是怎样展开的
论文中的 Blog 页面由两个独立区域组成:
Fetch<Page> blog = appCombine(
Page::new,
leftPane(),
mainPane());leftPane 又由 topics 和 popularPosts 组成。每个区域内部还会根据 post ID 查询浏览量、元数据和正文。
从源码结构看,这是一棵模块化的业务程序树;从执行时间看,它会变成几条按 round 推进的请求流水线:
Round 1
getPostIds
Round 2
getPostInfo(postIds...) ┐
getPostViews(postIds...) ┴─ 同一轮发现,按类型 batch
Round 3
getPostContent(postIds...)
以及依赖前两轮结果才确定的其他请求同一个 getPostIds 可能被不同模块引用。完整 Haxl 会用一次运行内的 request cache 让这些模块共享同一结果,因此模块不需要互相协调。缓存的意义不只是在命中时更快,而是允许每个模块独立声明自己需要的数据。
“最大并发”究竟是什么意思
更准确的表述是:Fetch 能在程序明确暴露的 Applicative 结构内,持续探索所有当前可达分支,直到遇到真实的数据依赖。于是:
- 同一依赖层的请求可以进入同一个 round;
- 同类请求可以 batch;
- 不同数据源的请求具备并发执行的机会;
- round 数接近查询依赖图的关键路径长度,而不是请求总数。
但这不等于任何实现都已经获得物理上的最大并行度。
当前公开的 Java Fetch 版本中,IO.parallel 仍然使用顺序 forEach(IO::performIO)。因此它已经实现了并发机会的发现和同轮调度,也实现了同类请求的 batch;但不同 request class 的 I/O 并没有真正在线程或异步 runtime 上同时执行。
同样,当前版本的 distinct() 只对一次 resolver 调用中的同类 ID 去重,并不是论文中贯穿一次 runFetch 的 request cache。完整 Haxl 还包含 typed fetch status、异常传播和独立 DataSource 等机制,这个 Java 版本没有全部移植。
这一区分很重要:
program exposes concurrency
≠
runtime physically executes I/O in parallel前者是 Fetch 最核心、也最难获得的结构;后者可以在 resolver 层继续演进。
真正消失的是手工协调
Fetch 的价值不只是把 94 次数据库请求压缩成 3 次。手工优化同样能做到这一点。
更重要的是,程序员可以继续按照业务结构拆分和复用查询:
- 一个模块声明自己需要什么数据;
- flatMap 保留必须等待的依赖;
- appCombine 和 mapM 暴露独立分支;
- 解释器不断探索到 Blocked;
- resolver 统一 batch 并执行请求;
- continuation 在结果回来后恢复。
于是性能优化不再要求所有模块共享一个巨大的编排函数,也不再要求调用者提前知道每个 Repository 的批量接口。
“There is no fork” 最终表达的是一种职责划分:业务程序描述数据依赖,Fetch runtime 决定怎样执行这些依赖。
这和 how fetch solve the problem 中“计划—优化—执行”的总结是同一件事。Fetch 把控制流变成可组合、可暂停、可检查的值,才获得了普通命令式查询代码没有的全局优化空间。
相关资料
- @There is no fork%3A an abstraction for efficient, concurrent, and concise data access