批处理
通过批处理请求并减少冗余 API 调用优化性能,提升数据获取与处理的效率。
在典型的应用开发中,当我们需要与外部 API、数据库或其他数据源交互时,常常会定义一些函数来发起请求,并相应地处理它们的结果或失败。
简单的模型搭建
下面是一个基础模型,它勾勒出我们的数据结构以及可能出现的错误:
import { Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
new GetUserError()._tag // => "GetUserError"
在真实场景中,我们可能希望为标识符使用更精确的类型,而不是直接使用原始类型(参见品牌类型)。 此外,你可能还想在错误中包含更详细的信息。
定义 API 函数
接下来我们定义一些与外部 API 交互的函数,处理诸如获取 Todo 列表、查询用户详情和发送邮件这样的常见操作。
import { Effect, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// API
// ------------------------------
// Fetches a list of todos from an external API
const getTodos = Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/todos").then(
(res) => res.json() as Promise<Array<Todo>>,
),
catch: () => new GetTodosError(),
})
// Retrieves a user by their ID from an external API
const getUserById = (id: number) =>
Effect.tryPromise({
try: () =>
fetch(`https://api.example.demo/getUserById?id=${id}`).then(
(res) => res.json() as Promise<User>,
),
catch: () => new GetUserError(),
})
// Sends an email via an external API
const sendEmail = (address: string, text: string) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/sendEmail", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({ address, text }),
}).then((res) => res.json() as Promise<void>),
catch: () => new SendEmailError(),
})
// Sends an email to a user by fetching their details first
const sendEmailToUser = (id: number, message: string) =>
getUserById(id).pipe(Effect.andThen((user) => sendEmail(user.email, message)))
// Notifies the owner of a todo by sending them an email
const notifyOwner = (todo: Todo) =>
getUserById(todo.ownerId).pipe(
Effect.andThen((user) =>
sendEmailToUser(user.id, `hey ${user.name} you got a todo!`),
),
)
new SendEmailError()._tag // => "SendEmailError"
在真实场景中,你不能一味相信 API 总会返回预期的数据 —— 为此,你可以使用
effect/Schema,或者 zod 之类的类似方案。
虽然这种做法直观易读,但未必最高效。重复的 API 调用,尤其是当多个 Todo 属于同一个所有者时,会显著增加网络开销,拖慢应用。
使用这些 API 函数
这些函数清晰易懂,但使用它们的方式未必最高效。例如,通知 Todo 所有者会涉及重复的 API 调用,而这部分是可以优化的。
import { Effect, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// API
// ------------------------------
// Fetches a list of todos from an external API
const getTodos = Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/todos").then(
(res) => res.json() as Promise<Array<Todo>>,
),
catch: () => new GetTodosError(),
})
// Retrieves a user by their ID from an external API
const getUserById = (id: number) =>
Effect.tryPromise({
try: () =>
fetch(`https://api.example.demo/getUserById?id=${id}`).then(
(res) => res.json() as Promise<User>,
),
catch: () => new GetUserError(),
})
// Sends an email via an external API
const sendEmail = (address: string, text: string) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/sendEmail", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({ address, text }),
}).then((res) => res.json() as Promise<void>),
catch: () => new SendEmailError(),
})
// Sends an email to a user by fetching their details first
const sendEmailToUser = (id: number, message: string) =>
getUserById(id).pipe(Effect.andThen((user) => sendEmail(user.email, message)))
// Notifies the owner of a todo by sending them an email
const notifyOwner = (todo: Todo) =>
getUserById(todo.ownerId).pipe(
Effect.andThen((user) =>
sendEmailToUser(user.id, `hey ${user.name} you got a todo!`),
),
)
// Orchestrates operations on todos, notifying their owners
const program = Effect.gen(function* () {
const todos = yield* getTodos
yield* Effect.forEach(todos, (todo) => notifyOwner(todo), {
concurrency: "unbounded",
})
})
new GetTodosError()._tag // => "GetTodosError"
这个实现会为每个 Todo 分别执行一次 API 调用,以获取所有者详情并发送邮件。如果多个 Todo 属于同一个所有者,就会产生冗余的 API 调用。
为了优化,如果你的后端支持批量 API 调用,可以考虑实现它。它把多个操作合并到单个请求中, 从而减少 HTTP 请求的数量,提升性能并降低负载。
批处理
假设 getUserById 和 sendEmail 可以批量执行。这意味着我们能在一次 HTTP 调用中发送多个请求,从而减少 API 请求数量并提升性能。
批处理的分步指南
-
声明请求: 我们首先把请求转换成结构化的数据模型。这需要详细描述输入参数、预期输出以及可能出现的错误。以这种方式组织请求,不仅有助于高效地管理数据,还能比较不同的请求,判断它们是否引用了相同的输入参数。
-
声明 Resolver: Resolver 旨在同时处理多个请求。借助比较请求的能力(确保它们引用相同的输入参数),Resolver 可以一次性执行多个请求,从而最大限度地发挥批处理的价值。
-
定义查询: 最后,我们定义一些查询,利用这些批量 Resolver 来执行操作。这一步把结构化的请求及其对应的 Resolver 组合成应用中可用的组成部分。
至关重要的是,请求的建模方式必须让它们可以相互比较。也就是说,要实现可比性(使用 Equals.equals 之类的方法),以便有效地识别并批量处理相同的请求。
声明请求
我们将借助 Request 这一概念,设计一个数据源可能支持的模型:
Request<Value, Error>
Request 是一种构造,表示对类型为 Value 的值的请求,它可能以类型为 Error 的错误失败。
我们先为数据源能够处理的各类请求定义一个结构化模型。
import { Request, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// Requests
// ------------------------------
// Define a request to get multiple Todo items which might
// fail with a GetTodosError
interface GetTodos extends Request.Request<Array<Todo>, GetTodosError> {
readonly _tag: "GetTodos"
}
// Create a tagged constructor for GetTodos requests
const GetTodos = Request.tagged<GetTodos>("GetTodos")
// Define a request to fetch a User by ID which might
// fail with a GetUserError
interface GetUserById extends Request.Request<User, GetUserError> {
readonly _tag: "GetUserById"
readonly id: number
}
// Create a tagged constructor for GetUserById requests
const GetUserById = Request.tagged<GetUserById>("GetUserById")
// Define a request to send an email which might
// fail with a SendEmailError
interface SendEmail extends Request.Request<void, SendEmailError> {
readonly _tag: "SendEmail"
readonly address: string
readonly text: string
}
// Create a tagged constructor for SendEmail requests
const SendEmail = Request.tagged<SendEmail>("SendEmail")
GetTodos()._tag // => "GetTodos"
每个请求都用一个具体的数据结构来定义,它继承自通用的 Request 类型,从而确保每个请求都携带自己特有的数据需求以及特定的错误类型。
通过使用 Request.tagged 这类带标签的构造器,我们可以轻松创建请求对象,使它们在整个应用中都能被识别和管理。
声明 Resolver
定义好请求之后,下一步是配置 Effect 如何使用 RequestResolver 解析这些请求:
RequestResolver<A>
RequestResolver<A> 能够执行类型为 A 的请求。传给 Effect.request 的 Resolver 自身不能带有任何依赖要求;关于如何在构造 Resolver 之前解析服务,请参见带上下文的 Resolver。
本节中,我们会为每种请求分别创建独立的 Resolver。Resolver 的粒度可以不同,但通常按照对应 API 调用是否支持批量处理来划分。
import { Effect, Request, RequestResolver, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// Requests
// ------------------------------
// Define a request to get multiple Todo items which might
// fail with a GetTodosError
interface GetTodos extends Request.Request<Array<Todo>, GetTodosError> {
readonly _tag: "GetTodos"
}
// Create a tagged constructor for GetTodos requests
const GetTodos = Request.tagged<GetTodos>("GetTodos")
// Define a request to fetch a User by ID which might
// fail with a GetUserError
interface GetUserById extends Request.Request<User, GetUserError> {
readonly _tag: "GetUserById"
readonly id: number
}
// Create a tagged constructor for GetUserById requests
const GetUserById = Request.tagged<GetUserById>("GetUserById")
// Define a request to send an email which might
// fail with a SendEmailError
interface SendEmail extends Request.Request<void, SendEmailError> {
readonly _tag: "SendEmail"
readonly address: string
readonly text: string
}
// Create a tagged constructor for SendEmail requests
const SendEmail = Request.tagged<SendEmail>("SendEmail")
// ------------------------------
// Resolvers
// ------------------------------
// Assuming GetTodos cannot be batched, we create a standard resolver
const GetTodosResolver = RequestResolver.fromEffect(
(_: Request.Entry<GetTodos>): Effect.Effect<Todo[], GetTodosError> =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/todos").then(
(res) => res.json() as Promise<Array<Todo>>,
),
catch: () => new GetTodosError(),
}),
)
// Assuming GetUserById can be batched, we create a batched resolver
const GetUserByIdResolver = RequestResolver.make(
(entries: ReadonlyArray<Request.Entry<GetUserById>>) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/getUserByIdBatch", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({
users: entries.map(({ request }) => ({ id: request.id })),
}),
}).then((res) => res.json()) as Promise<Array<User>>,
catch: () => new GetUserError(),
}).pipe(
Effect.andThen((users) =>
Effect.forEach(entries, (entry, index) =>
Request.completeEffect(entry, Effect.succeed(users[index]!)),
),
),
Effect.catch((error) =>
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.fail(error)),
),
),
),
)
// Assuming SendEmail can be batched, we create a batched resolver
const SendEmailResolver = RequestResolver.make(
(entries: ReadonlyArray<Request.Entry<SendEmail>>) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/sendEmailBatch", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({
emails: entries.map(({ request }) => ({
address: request.address,
text: request.text,
})),
}),
}).then((res) => res.json() as Promise<void>),
catch: () => new SendEmailError(),
}).pipe(
Effect.andThen(
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.void),
),
),
Effect.catch((error) =>
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.fail(error)),
),
),
),
)
SendEmail({ address: "a@b.com", text: "hi" })._tag // => "SendEmail"
与其他 effect 一样,Resolver 也可以访问上下文,而且创建 Resolver 的方式多种多样。 更多细节请参阅 RequestResolver 模块的参考文档。
在这个配置中:
- GetTodosResolver 负责获取多个
Todo项。因为我们假设它不能批量处理,所以把它配置为普通 Resolver。 - GetUserByIdResolver 和 SendEmailResolver 被配置为批量 Resolver。这样设置的前提是这些请求可以按批处理,从而提升性能并减少 API 调用次数。
定义查询
现在解析器已经就绪,我们可以把所有部分串联起来定义查询了。这一步让我们能够在应用中高效地执行数据操作。
import { Effect, Request, RequestResolver, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// Requests
// ------------------------------
// Define a request to get multiple Todo items which might
// fail with a GetTodosError
interface GetTodos extends Request.Request<Array<Todo>, GetTodosError> {
readonly _tag: "GetTodos"
}
// Create a tagged constructor for GetTodos requests
const GetTodos = Request.tagged<GetTodos>("GetTodos")
// Define a request to fetch a User by ID which might
// fail with a GetUserError
interface GetUserById extends Request.Request<User, GetUserError> {
readonly _tag: "GetUserById"
readonly id: number
}
// Create a tagged constructor for GetUserById requests
const GetUserById = Request.tagged<GetUserById>("GetUserById")
// Define a request to send an email which might
// fail with a SendEmailError
interface SendEmail extends Request.Request<void, SendEmailError> {
readonly _tag: "SendEmail"
readonly address: string
readonly text: string
}
// Create a tagged constructor for SendEmail requests
const SendEmail = Request.tagged<SendEmail>("SendEmail")
// ------------------------------
// Resolvers
// ------------------------------
// Assuming GetTodos cannot be batched, we create a standard resolver
const GetTodosResolver = RequestResolver.fromEffect(
(_: Request.Entry<GetTodos>): Effect.Effect<Todo[], GetTodosError> =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/todos").then(
(res) => res.json() as Promise<Array<Todo>>,
),
catch: () => new GetTodosError(),
}),
)
// Assuming GetUserById can be batched, we create a batched resolver
const GetUserByIdResolver = RequestResolver.make(
(entries: ReadonlyArray<Request.Entry<GetUserById>>) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/getUserByIdBatch", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({
users: entries.map(({ request }) => ({ id: request.id })),
}),
}).then((res) => res.json()) as Promise<Array<User>>,
catch: () => new GetUserError(),
}).pipe(
Effect.andThen((users) =>
Effect.forEach(entries, (entry, index) =>
Request.completeEffect(entry, Effect.succeed(users[index]!)),
),
),
Effect.catch((error) =>
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.fail(error)),
),
),
),
)
// Assuming SendEmail can be batched, we create a batched resolver
const SendEmailResolver = RequestResolver.make(
(entries: ReadonlyArray<Request.Entry<SendEmail>>) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/sendEmailBatch", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({
emails: entries.map(({ request }) => ({
address: request.address,
text: request.text,
})),
}),
}).then((res) => res.json() as Promise<void>),
catch: () => new SendEmailError(),
}).pipe(
Effect.andThen(
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.void),
),
),
Effect.catch((error) =>
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.fail(error)),
),
),
),
)
// ------------------------------
// Queries
// ------------------------------
// Defines a query to fetch all Todo items
const getTodos: Effect.Effect<Array<Todo>, GetTodosError> = Effect.request(
GetTodos(),
GetTodosResolver,
)
// Defines a query to fetch a user by their ID
const getUserById = (id: number) =>
Effect.request(GetUserById({ id }), GetUserByIdResolver)
// Defines a query to send an email to a specific address
const sendEmail = (address: string, text: string) =>
Effect.request(SendEmail({ address, text }), SendEmailResolver)
// Composes getUserById and sendEmail to send an email to a specific user
const sendEmailToUser = (id: number, message: string) =>
getUserById(id).pipe(Effect.andThen((user) => sendEmail(user.email, message)))
// Uses getUserById to fetch the owner of a Todo and then sends them an email notification
const notifyOwner = (todo: Todo) =>
getUserById(todo.ownerId).pipe(
Effect.andThen((user) =>
sendEmailToUser(user.id, `hey ${user.name} you got a todo!`),
),
)
GetUserById({ id: 1 }).id // => 1
通过使用 Effect.request 函数,我们让解析器与请求模型有效地结合在一起。这种方式确保每个查询都能用恰当的解析器以最优方式完成。
尽管代码结构与前面的示例看起来相似,但使用解析器能显著提升效率:它优化了请求的处理方式,并减少了不必要的 API 调用。
const program = Effect.gen(function* () {
const todos = yield* getTodos
yield* Effect.forEach(todos, (todo) => notifyOwner(todo), {
batching: true,
})
})
在最终的配置下,无论有多少个 todo,这个程序都只会向 API 执行 3 次查询。这与传统方式形成鲜明对比:后者可能执行 1 + 2n 次查询,其中 n 是 todo 的数量。这是效率上的显著提升,尤其是对于数据交互量很大的应用而言。
带上下文的解析器
在复杂的应用中,解析器通常需要访问共享服务或配置,才能有效地处理请求。然而,在提供必要上下文的同时保持请求批处理能力,可能颇具挑战。这里我们将探讨如何在解析器中管理上下文,以确保批处理能力不受影响。
在创建请求解析器时,谨慎管理上下文至关重要。为解析器提供过多的上下文,或给不同的解析器提供不同的服务,都会使它们无法兼容批处理。为避免这类问题,传给 Effect.request 的解析器,其上下文被显式设为 never。这迫使开发者明确界定上下文在解析器内部是如何被访问和使用的。
考虑下面的例子,我们搭建了一个 HTTP 服务,供解析器用来执行 API 调用:
import { Effect, Context, RequestResolver, Request, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// Requests
// ------------------------------
// Define a request to get multiple Todo items which might
// fail with a GetTodosError
interface GetTodos extends Request.Request<
Array<Todo>,
GetTodosError,
HttpService
> {
readonly _tag: "GetTodos"
}
// Create a tagged constructor for GetTodos requests
const GetTodos = Request.tagged<GetTodos>("GetTodos")
// Define a request to fetch a User by ID which might
// fail with a GetUserError
interface GetUserById extends Request.Request<User, GetUserError> {
readonly _tag: "GetUserById"
readonly id: number
}
// Create a tagged constructor for GetUserById requests
const GetUserById = Request.tagged<GetUserById>("GetUserById")
// Define a request to send an email which might
// fail with a SendEmailError
interface SendEmail extends Request.Request<void, SendEmailError> {
readonly _tag: "SendEmail"
readonly address: string
readonly text: string
}
// Create a tagged constructor for SendEmail requests
const SendEmail = Request.tagged<SendEmail>("SendEmail")
// ------------------------------
// Resolvers With Context
// ------------------------------
class HttpService extends Context.Service<
HttpService,
{ fetch: typeof fetch }
>()("HttpService") {}
// the resolver itself must have no requirements, so we resolve HttpService
// before constructing it, and let that outer effect carry the requirement
const GetTodosResolver = Effect.map(HttpService, (http) =>
RequestResolver.fromEffect(
(_: Request.Entry<GetTodos>): Effect.Effect<Array<Todo>, GetTodosError> =>
Effect.tryPromise({
try: () =>
http
.fetch("https://api.example.demo/todos")
.then((res) => res.json() as Promise<Array<Todo>>),
catch: () => new GetTodosError(),
}),
),
)
HttpService.key // => "HttpService"
现在可以看到,GetTodosResolver 的类型不再是 RequestResolver,而是:
const GetTodosResolver: Effect<RequestResolver<GetTodos>, never, HttpService>
这是一个 effect,它访问 HttpService,并返回一个已经组装好、具备最小可用上下文的解析器。
有了这样一个 effect,我们就可以直接在查询定义中使用它:
const getTodos: Effect.Effect<Todo[], GetTodosError, HttpService> =
Effect.request(GetTodos(), GetTodosResolver)
可以看到,这个 Effect 正确地要求提供 HttpService。
另一种做法是,把 RequestResolver 作为 Layer 的一部分来创建,在构造时直接访问上下文,或闭包捕获上下文。
示例
import { Effect, Context, RequestResolver, Request, Layer, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// Requests
// ------------------------------
// Define a request to get multiple Todo items which might
// fail with a GetTodosError
interface GetTodos extends Request.Request<Array<Todo>, GetTodosError> {
readonly _tag: "GetTodos"
}
// Create a tagged constructor for GetTodos requests
const GetTodos = Request.tagged<GetTodos>("GetTodos")
// Define a request to fetch a User by ID which might
// fail with a GetUserError
interface GetUserById extends Request.Request<User, GetUserError> {
readonly _tag: "GetUserById"
readonly id: number
}
// Create a tagged constructor for GetUserById requests
const GetUserById = Request.tagged<GetUserById>("GetUserById")
// Define a request to send an email which might
// fail with a SendEmailError
interface SendEmail extends Request.Request<void, SendEmailError> {
readonly _tag: "SendEmail"
readonly address: string
readonly text: string
}
// Create a tagged constructor for SendEmail requests
const SendEmail = Request.tagged<SendEmail>("SendEmail")
// ------------------------------
// Resolvers With Context
// ------------------------------
class HttpService extends Context.Service<
HttpService,
{ fetch: typeof fetch }
>()("HttpService") {}
// the resolver itself must have no requirements, so we resolve HttpService
// before constructing it, and let that outer effect carry the requirement
const GetTodosResolver = Effect.map(HttpService, (http) =>
RequestResolver.fromEffect(
(_: Request.Entry<GetTodos>): Effect.Effect<Array<Todo>, GetTodosError> =>
Effect.tryPromise({
try: () =>
http
.fetch("https://api.example.demo/todos")
.then((res) => res.json() as Promise<Array<Todo>>),
catch: () => new GetTodosError(),
}),
),
)
// ------------------------------
// Layers
// ------------------------------
class TodosService extends Context.Service<
TodosService,
{
getTodos: Effect.Effect<Array<Todo>, GetTodosError>
}
>()("TodosService") {}
const TodosServiceLive = Layer.effect(
TodosService,
Effect.gen(function* () {
const http = yield* HttpService
const resolver = RequestResolver.fromEffect((_: Request.Entry<GetTodos>) =>
Effect.tryPromise({
try: () =>
http
.fetch("https://api.example.demo/todos")
.then<any, Todo[]>((res) => res.json()),
catch: () => new GetTodosError(),
}),
)
return {
getTodos: Effect.request(GetTodos(), resolver),
}
}),
)
const getTodos: Effect.Effect<
Array<Todo>,
GetTodosError,
TodosService
> = Effect.andThen(TodosService, (service) => service.getTodos)
TodosService.key // => "TodosService"
鉴于 Layer 是把服务装配到一起的自然原语,对大多数场景而言,这种方式很可能也是最好的。
缓存
虽然我们已经大幅优化了请求批处理,但还有一个领域可以进一步提升应用的效率:缓存。没有缓存时,即使批处理已经过优化,相同的请求仍可能被执行多次,导致不必要的数据获取。
缓存配置在 resolver 上。RequestResolver.withCache 会把 resolver 包装进一个以请求相等性为键的有界内存缓存。任何基于该 resolver、用 Effect.request 构建的查询都会自动沿用相同的缓存行为。
下面是为 getUserById 查询实现缓存的方式:
import { Effect, Request, RequestResolver, Data } from "effect"
// ------------------------------
// Model
// ------------------------------
interface User {
readonly _tag: "User"
readonly id: number
readonly name: string
readonly email: string
}
class GetUserError extends Data.TaggedError("GetUserError")<{}> {}
interface Todo {
readonly _tag: "Todo"
readonly id: number
readonly message: string
readonly ownerId: number
}
class GetTodosError extends Data.TaggedError("GetTodosError")<{}> {}
class SendEmailError extends Data.TaggedError("SendEmailError")<{}> {}
// ------------------------------
// Requests
// ------------------------------
// Define a request to get multiple Todo items which might
// fail with a GetTodosError
interface GetTodos extends Request.Request<Array<Todo>, GetTodosError> {
readonly _tag: "GetTodos"
}
// Create a tagged constructor for GetTodos requests
const GetTodos = Request.tagged<GetTodos>("GetTodos")
// Define a request to fetch a User by ID which might
// fail with a GetUserError
interface GetUserById extends Request.Request<User, GetUserError> {
readonly _tag: "GetUserById"
readonly id: number
}
// Create a tagged constructor for GetUserById requests
const GetUserById = Request.tagged<GetUserById>("GetUserById")
// Define a request to send an email which might
// fail with a SendEmailError
interface SendEmail extends Request.Request<void, SendEmailError> {
readonly _tag: "SendEmail"
readonly address: string
readonly text: string
}
// Create a tagged constructor for SendEmail requests
const SendEmail = Request.tagged<SendEmail>("SendEmail")
// ------------------------------
// Resolvers
// ------------------------------
// Assuming GetTodos cannot be batched, we create a standard resolver
const GetTodosResolver = RequestResolver.fromEffect(
(_: Request.Entry<GetTodos>): Effect.Effect<Todo[], GetTodosError> =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/todos").then(
(res) => res.json() as Promise<Array<Todo>>,
),
catch: () => new GetTodosError(),
}),
)
// Assuming GetUserById can be batched, we create a batched resolver
const GetUserByIdResolver = RequestResolver.make(
(entries: ReadonlyArray<Request.Entry<GetUserById>>) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/getUserByIdBatch", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({
users: entries.map(({ request }) => ({ id: request.id })),
}),
}).then((res) => res.json()) as Promise<Array<User>>,
catch: () => new GetUserError(),
}).pipe(
Effect.andThen((users) =>
Effect.forEach(entries, (entry, index) =>
Request.completeEffect(entry, Effect.succeed(users[index]!)),
),
),
Effect.catch((error) =>
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.fail(error)),
),
),
),
)
// Assuming SendEmail can be batched, we create a batched resolver
const SendEmailResolver = RequestResolver.make(
(entries: ReadonlyArray<Request.Entry<SendEmail>>) =>
Effect.tryPromise({
try: () =>
fetch("https://api.example.demo/sendEmailBatch", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({
emails: entries.map(({ request }) => ({
address: request.address,
text: request.text,
})),
}),
}).then((res) => res.json() as Promise<void>),
catch: () => new SendEmailError(),
}).pipe(
Effect.andThen(
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.void),
),
),
Effect.catch((error) =>
Effect.forEach(entries, (entry) =>
Request.completeEffect(entry, Effect.fail(error)),
),
),
),
)
// ------------------------------
// Caching
// ------------------------------
// Wrap the resolver in a bounded, in-memory cache keyed by request equality.
// Build it once and reuse the resulting resolver for every call.
const cachedGetUserByIdResolver = Effect.runSync(
RequestResolver.withCache(GetUserByIdResolver, { capacity: 256 }),
)
const getUserById = (id: number) =>
Effect.request(GetUserById({ id }), cachedGetUserByIdResolver)
cachedGetUserByIdResolver === GetUserByIdResolver // => false
最终程序
假设你已经把所有部分正确串联起来:
const program = Effect.gen(function* () {
const todos = yield* getTodos
yield* Effect.forEach(todos, (todo) => notifyOwner(todo), {
concurrency: "unbounded",
})
}).pipe(Effect.repeat(Schedule.fixed("10 seconds")))
在这个程序中,getTodos 操作会获取每个用户的 todo。随后,Effect.forEach 函数用于并发地通知每个 todo 的所有者,而无需等待这些通知完成。
repeat 函数被应用于整条操作链,它使用固定调度(fixed schedule)确保程序每 10 秒重复一次。这意味着整个流程——包括获取 todo 与发送通知——都会以 10 秒为间隔反复执行。
由于 getUserById 建立在 cachedGetUserByIdResolver 之上,只要该 resolver 的缓存尚未淘汰对应的 GetUserById 请求(淘汰受传给 RequestResolver.withCache 的 capacity 限制),程序就会复用该请求的缓存结果,从而减少获取用户数据的不必要请求。
此外,该程序设计为批量发送邮件,从而实现高效处理并更好地利用资源。
自定义请求缓存
RequestResolver.withCache 接受一个 strategy 选项("lru"(默认)或 "fifo"),用于控制当缓存达到 capacity 后淘汰哪个条目:
const fifoCachedResolver = Effect.runSync(
RequestResolver.withCache(GetUserByIdResolver, {
capacity: 256,
strategy: "fifo",
}),
)
以这种方式创建的缓存条目不会按时间过期。如果你需要基于存活时间(time-to-live)的过期机制,或者希望把缓存查找暴露为一等的 Cache(带有 get/refresh/invalidate,参见 Cache)而不是普通的 RequestResolver,请改用 RequestResolver.asCache。它接受同样的 capacity 选项,外加一个可选的 timeToLive。