已发布 上游基线 bf46254 原文 ↗ 在 GitHub 编辑

批处理

通过批处理请求并减少冗余 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"
Use Precise Types and Detailed Errors

在真实场景中,我们可能希望为标识符使用更精确的类型,而不是直接使用原始类型(参见品牌类型)。 此外,你可能还想在错误中包含更详细的信息。

定义 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"
Validating API Responses

在真实场景中,你不能一味相信 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 调用。

Improving Efficiency with Batch Calls

为了优化,如果你的后端支持批量 API 调用,可以考虑实现它。它把多个操作合并到单个请求中, 从而减少 HTTP 请求的数量,提升性能并降低负载。

批处理

假设 getUserByIdsendEmail 可以批量执行。这意味着我们能在一次 HTTP 调用中发送多个请求,从而减少 API 请求数量并提升性能。

批处理的分步指南

  1. 声明请求: 我们首先把请求转换成结构化的数据模型。这需要详细描述输入参数、预期输出以及可能出现的错误。以这种方式组织请求,不仅有助于高效地管理数据,还能比较不同的请求,判断它们是否引用了相同的输入参数。

  2. 声明 Resolver: Resolver 旨在同时处理多个请求。借助比较请求的能力(确保它们引用相同的输入参数),Resolver 可以一次性执行多个请求,从而最大限度地发挥批处理的价值。

  3. 定义查询: 最后,我们定义一些查询,利用这些批量 Resolver 来执行操作。这一步把结构化的请求及其对应的 Resolver 组合成应用中可用的组成部分。

Ensuring Request Comparability

至关重要的是,请求的建模方式必须让它们可以相互比较。也就是说,要实现可比性(使用 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"
Accessing Context in Resolvers

与其他 effect 一样,Resolver 也可以访问上下文,而且创建 Resolver 的方式多种多样。 更多细节请参阅 RequestResolver 模块的参考文档。

在这个配置中:

  • GetTodosResolver 负责获取多个 Todo 项。因为我们假设它不能批量处理,所以把它配置为普通 Resolver。
  • GetUserByIdResolverSendEmailResolver 被配置为批量 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.withCachecapacity 限制),程序就会复用该请求的缓存结果,从而减少获取用户数据的不必要请求。

此外,该程序设计为批量发送邮件,从而实现高效处理并更好地利用资源。

自定义请求缓存

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