扩展

Drain 流水线

该流水线包装每个 drain:它会批量处理事件、在失败时重试、将事件分发到多个目标,并发送浏览器日志。

在生产环境中,为每个发出的事件发送一个 HTTP 请求无法扩展。drain pipeline 会缓冲事件并批量发送,在临时失败时重试,缓冲区溢出时丢弃最旧的事件,并允许单个 drain 函数并行分发到多个目标。用于客户端日志的 HTTP browser drain 也由同一条流水线提供支持。

drain pipeline · batch + retry·BUFFERING
app emits#0
never blocks the response
in-memory buffer0 / 10
·
·
·
·
·
·
·
·
·
·
fills until size or 5s interval
flush · http POST0 ok
onDropped · maxBufferSize · drain.flush()
emitted0
batches sent0
dropped0
你想要…查看
将任意 drain 包装为批量 + 重试 + 缓冲快速开始
将每个事件并行发送到多个目标Fanout
将浏览器日志发送到你的服务器端点HTTP drain(浏览器到服务器)
调整批量大小、重试策略、缓冲区大小配置

添加 drain 流水线(批量 + 重试 + 分发)

快速开始

该流水线可以包装任意 drain。具体接入方式取决于你的框架,因此请选择与你使用的框架匹配的选项卡;下面的其他示例都使用相同的结构。

// server/plugins/evlog-drain.ts
import type { DrainContext } from 'evlog'
import { createDrainPipeline } from 'evlog/pipeline'
import { createAxiomDrain } from 'evlog/axiom'

export default defineNitroPlugin((nitroApp) => {
  const pipeline = createDrainPipeline<DrainContext>()
  const drain = pipeline(createAxiomDrain())

  nitroApp.hooks.hook('evlog:drain', drain)
  nitroApp.hooks.hook('close', () => drain.flush())
})
始终在进程退出前刷新流水线(drain.flush())。在 Nitro 中使用 close hook;在 standalone 脚本中于 process.exit 前调用。在 serverless 运行时中,中间件会为你延长实例的生命周期,参见 Serverless

流水线对事件的处理方式

事件到达 evlog:drain 后会被缓冲。当达到 batch.sizebatch.intervalMs 到期时(以先发生者为准),批次会被刷新。失败时,相同批次会按照配置的退避策略重试;一旦 retry.maxAttempts 耗尽,就会使用丢失的事件调用 onDropped(如果没有配置该回调,则记录丢弃日志)。缓冲区大小受 maxBufferSize 限制。缓冲区填满后,会丢弃最旧的事件,以保持内存占用稳定。

流水线负责重试

直接调用时,Adapter drains(createAxiomDrain()createDatadogDrain() 等)会吞掉传输失败,因此失败的后端永远不会中断请求。在流水线内部,这会使 retryonDropped 成为无效代码,因此流水线会将适配器解包为会抛出错误的 raw 版本:失败会进入重试循环,而原始传输默认只进行一次 HTTP 尝试,以避免批次被重复重试。如果你明确希望在每次流水线尝试中额外进行 HTTP 尝试,请仅在适配器配置中设置 retries

配置

pipeline-config.ts
import type { DrainContext } from 'evlog'
import { createDrainPipeline } from 'evlog/pipeline'
import { createAxiomDrain } from 'evlog/axiom'

const pipeline = createDrainPipeline<DrainContext>({
  batch: {
    size: 50,          // 每 50 个事件刷新一次
    intervalMs: 5000,  // 或每 5 秒刷新一次,以先到者为准
  },
  retry: {
    maxAttempts: 3,
    backoff: 'exponential',
    initialDelayMs: 1000,
    maxDelayMs: 30000,
  },
  maxBufferSize: 1000,
  onDropped: (events, error) => {
    console.error(`[evlog] 丢弃了 ${events.length} 个事件:`, error?.message)
  },
})

export const drain = pipeline(createAxiomDrain())

每个选项及其默认值

选项默认值描述
batch.size50每批最多事件数
batch.intervalMs5000刷新部分批次前的最大时间(ms)
retry.maxAttempts3包括初始尝试在内的总尝试次数
retry.backoff'exponential''exponential' | 'linear' | 'fixed'
retry.initialDelayMs1000第一次重试的基础延迟
retry.maxDelayMs30000任意重试延迟的上限
maxBufferSize1000丢弃最旧事件前的最大缓冲事件数
onDropped-事件被丢弃时的回调(溢出或重试耗尽)

选择重试间隔策略

策略延迟模式使用场景
exponential1s, 2s, 4s, 8s…默认值。最适合可能需要时间恢复的临时失败
linear1s, 2s, 3s, 4s…可预测的延迟增长
fixed1s, 1s, 1s, 1s…每次延迟都相同。适合限流 API

包装器返回的内容

pipeline(drain) 返回的函数与 hook 兼容,并提供:

属性类型描述
drain(ctx)(ctx: T) => void将单个事件推入缓冲区
drain.flush()() => Promise<void>强制刷新所有缓冲的事件
drain.settled()() => Promise<void>当目前为止缓冲的所有内容都已发送或丢弃后解析,不会强制刷新
drain.pendingnumber当前缓冲的事件数

Serverless

drain(ctx) 会在事件进入缓冲区后返回,因此在 serverless 运行时中,关键问题是:在批次实际发送完成前,是什么让实例保持运行。当你将流水线作为框架集成的 drain 选项传递时,中间件会向运行时的 waitUntil 注册 drain.settled()(Next 的 after()、Workers 的 ctx.waitUntil):实例会保持运行,直到缓冲的批次发送或丢弃。batch.intervalMs 会限制首次传输尝试开始的时间;重试退避可能会使生命周期超过该时间,直到传输成功或 retry.maxAttempts 耗尽。如果你希望缩短响应后的生命周期,而不是优先进行批处理,请在 serverless 环境中降低 intervalMs

Nitro 是例外:在 evlog:drain hook 上注册的 drains 对插件来说是不透明的,因此插件无法为你延长生命周期。在 Node 中,快速开始示例中的 close hook 刷新操作可以覆盖关闭流程;在 Cloudflare 中,请在处理请求的地方自行向执行上下文的 waitUntil 注册 drain.settled()

Fanout

drain pipeline
0/5 delivered
wide event
{
  level:    "info",
  method:   "POST",
  path:     "/checkout",
  duration: 234,
  user:     {...},
  cart:     {...}
}
BUILDING
  1. axiom
    pending
  2. otlp
    pending
  3. sentry
    pending
  4. nuxthub
    pending
  5. fs
    pending
fire-and-forget retry with backoff one failure ≠ blocked response ✓ all destinations resolved

在生产环境中,同一类宽泛事件通常需要到达多个目标:长期存储(Axiom)、指标工具(Datadog)、错误追踪器(Sentry),以及用于事件回放的本地 fs drain。流水线先批量处理一次,然后你的 drain 通过 Promise.allSettled 将该批次并行分发到每个目标,这样单个缓慢/失败的目标就不会阻塞其他目标。

将 evlog 事件分发到多个目标

接入 fan-out

import { createDrainPipeline } from 'evlog/pipeline'
import { createAxiomDrain } from 'evlog/axiom'
import { createDatadogDrain } from 'evlog/datadog'
import { createSentryDrain } from 'evlog/sentry'
import { createFsDrain } from 'evlog/fs'
import type { DrainContext } from 'evlog'

const pipeline = createDrainPipeline<DrainContext>({
  batch: { size: 50, intervalMs: 5000 },
  retry: { maxAttempts: 3 },
  maxBufferSize: 1000,
})

const axiom = createAxiomDrain()
const datadog = createDatadogDrain()
const sentry = createSentryDrain({ minLevel: 'error' })
const fs = createFsDrain({ dir: '.evlog/logs', maxFiles: 14 })

export const drain = pipeline(async (batch) => {
  await Promise.allSettled([
    axiom(batch),
    datadog(batch),
    sentry(batch),
    fs(batch),
  ])
})

fan-out 的保证

  • 并行分发:每个目标都会通过 Promise.allSettled 并发接收该批次
  • 容错 fanout:如果 Datadog 的 API 抛出错误,Axiom / Sentry / fs 仍会完成;只有包装函数被拒绝时,流水线才会重试整个批次
  • 共享背压:整个流水线只设置一次缓冲区大小;如果包装的 drain 处理落后,所有目标都会一致地丢弃最旧事件

将不同事件发送到每个 drain

包装某个目标 drain,让它只看到你关心的事件:

import type { DrainContext } from 'evlog'

const sentry = createSentryDrain({ dsn: process.env.SENTRY_DSN! })

async function sentryErrorsOnly(batch: DrainContext[]): Promise<void> {
  const errors = batch.filter(c => c.event?.level === 'error')
  if (errors.length > 0) await sentry(errors)
}

export const drain = pipeline(async (batch) => {
  await Promise.allSettled([
    axiom(batch),
    sentryErrorsOnly(batch),
  ])
})

大多数内置 drains 都直接暴露 minLevel,因此你只在需要非 level 过滤(路径、自定义字段等)时才需要这种模式。

自己编写 drain

你不需要适配器。传入任意接受一个 batch 的异步函数:

pipeline-custom.ts
import type { DrainContext } from 'evlog'
import { createDrainPipeline } from 'evlog/pipeline'

const pipeline = createDrainPipeline<DrainContext>({ batch: { size: 100 } })

export const drain = pipeline(async (batch) => {
  await fetch('https://your-service.com/logs', {
    method: 'POST',
    headers: { 'Content-Type': 'application/json' },
    body: JSON.stringify(batch.map(ctx => ctx.event)),
  })
})

对于更复杂的场景(配置解析、重试、身份标头),请改用 defineHttpDrain,让工具包处理样板代码。

HTTP drain(浏览器到服务器)

HTTP drain 会将客户端日志发送到你的服务器,无论服务器使用什么框架。它构建于相同的流水线基础能力之上,并使用浏览器专用的默认设置(fetch keepalive + 在 visibilitychange 时使用 sendBeacon)。

evlog/browser 导入路径已被 弃用,它会重新导出与 evlog/http 相同的 API。它将在下一个 主版本 中移除。新代码请优先使用 evlog/http

为客户端日志设置 HTTP 传输

快速开始

app.ts
import { initLogger, log } from 'evlog'
import { createHttpLogDrain } from 'evlog/http'

const drain = createHttpLogDrain({
  drain: { endpoint: 'https://logs.example.com/v1/ingest' },
})
initLogger({ drain })

log.info({ action: 'page_view', path: location.pathname })

工作原理(浏览器特定)

  1. log.info() / log.warn() / log.error() 将事件推入内存缓冲区
  2. 事件会按大小(默认 25)或时间间隔(默认 2 秒)进行批处理
  3. 批次通过带有 keepalive: truefetch 发送,以确保请求在页面导航时仍能继续
  4. 当页面变为隐藏状态(切换标签、导航)时,缓冲的事件会通过 navigator.sendBeacon 作为回退方式刷新发送
  5. 你的服务器端点接收一个 DrainContext[] JSON 数组,并以你喜欢的方式处理它

选择所需层级

createHttpLogDrain(options)

高级、预组合:创建一个带有批处理、重试以及在 visibilitychange 时自动刷新的 pipeline。直接返回可与 initLogger({ drain }) 一起使用的 PipelineDrainFn<DrainContext>

app.ts
import { initLogger, log } from 'evlog'
import { createHttpLogDrain } from 'evlog/http'

const drain = createHttpLogDrain({
  drain: { endpoint: 'https://logs.example.com/v1/ingest' },
  pipeline: { batch: { size: 50, intervalMs: 5000 } },
})

initLogger({ drain })
log.info({ action: 'click', target: 'buy-button' })

createHttpDrain(config)

低层传输函数。当你希望完全控制 pipeline 配置时使用它:

app.ts
import { createHttpDrain } from 'evlog/http'
import { createDrainPipeline } from 'evlog/pipeline'
import type { DrainContext } from 'evlog'

const transport = createHttpDrain({
  endpoint: 'https://logs.example.com/v1/ingest',
})
const pipeline = createDrainPipeline<DrainContext>({
  batch: { size: 100, intervalMs: 10000 },
  retry: { maxAttempts: 5 },
})

const drain = pipeline(transport)

每个浏览器选项

HttpDrainConfig

选项默认值说明
endpoint-(必需) 服务器 ingest 端点的完整 URL
headers-随每个 fetch 请求发送的自定义标头(例如 AuthorizationX-API-Key
timeout5000请求超时时间,单位为毫秒
useBeacontrue当页面隐藏时使用 sendBeacon
credentials'same-origin'Fetch 凭据模式('omit''same-origin''include')。对于跨域端点请设置为 'include'

HttpLogDrainOptions

选项默认值说明
drain-(必需) HttpDrainConfig 对象
pipeline{ batch: { size: 25, intervalMs: 2000 }, retry: { maxAttempts: 2 } }pipeline 配置覆盖项
autoFlushtrue自动注册 visibilitychange 刷新监听器

让日志经受住标签页关闭

当启用 useBeacon(默认值)且页面变为隐藏时,drain 会自动从 fetch 切换到 navigator.sendBeacon。这可确保即使用户关闭标签页或离开页面,日志也能被发送。

sendBeacon 受浏览器施加的负载限制(约 64 KB)。如果负载超过此限制,drain 会抛出错误。请保持批次大小合理(默认 25 远在限制范围内)。

验证浏览器请求

传递自定义标头以保护你的 ingest 端点:

app.ts
const drain = createHttpLogDrain({
  drain: {
    endpoint: 'https://logs.example.com/v1/ingest',
    headers: {
      'Authorization': 'Bearer ' + token,
    },
  },
})
headers 仅应用于 fetch 请求。sendBeacon API 不支持自定义标头,因此当页面隐藏并使用 sendBeacon 时,不会发送标头。如果你的端点需要身份验证,请通过会话 cookie 进行验证(对跨域端点设置 credentials: 'include'),或者通过 useBeacon: false 禁用 sendBeacon

服务器端点

你的服务器需要一个接受 DrainContext[] JSON body 的 POST 端点。常见框架示例如下:

app.post('/v1/ingest', express.json(), (req, res) => {
  for (const entry of req.body) {
    console.log('[BROWSER]', JSON.stringify(entry))
  }
  res.sendStatus(204)
})

请参阅完整的 browser example,了解可运行的 Hono server + browser page。

会导致事件丢失的错误做法

  • 不要忘记在关闭时调用 drain.flush():否则缓冲的事件会丢失
  • 调整 batch.size 以匹配你的服务商建议的负载大小:太小会浪费开销,太大则有被拒绝的风险
  • 除非需要按目标分别重试,否则不要为每个 drain 运行一个流水线:共享一个流水线可以让批处理 + 缓冲保持一致
  • 在分发到与 serverless 不兼容的目标前先进行检查stream server 通过进程内 stream 将事件发送给每个已连接的客户端;它不是 drain

下一步