扩展
为任何没有内置适配器的后端构建 Drain。配置解析、重试、超时和身份标头都由你处理。

drain 是 evlog 流水线的终端步骤:一个接收宽事件并将其发送到某处的函数:HTTP API、消息队列、数据库、webhook、本地文件。evlog 为常见供应商提供了内置 drain(适配器概览)。当你需要发送到尚未覆盖的目标时,可以自行编写。

两个工厂覆盖所有情况:

你有……使用
一个 HTTP 后端(REST、JSON ingest、供应商 /v1/logs 端点)defineHttpDrain
一个非 HTTP 传输(gRPC、WebSocket、供应商 SDK、队列、原始套接字)defineDrain

两者都来自 evlog/toolkit,并且正是每个内置适配器所使用的工厂。

构建一个自定义 evlog drain

就是这样。defineHttpDrain 会处理重试(默认 2 次)、超时(默认 5000ms)、错误隔离以及身份标头(User-Agent: evlog/<version> + X-Evlog-Source: <name>)。即使目标服务宕机,你的应用流水线仍会继续运行;而当你将 drain 包装在 createDrainPipeline 中时,流水线会通过 drain 的 raw 变体接管重试。

一个 5 分钟示例:内部 Loki drain

一个完整可工作的 drain,只有 25 行,并且不需要外部配置助手:

lib/loki-drain.ts
import { defineHttpDrain } from 'evlog/toolkit'

export function createLokiDrain(overrides?: { url?: string, token?: string }) {
  return defineHttpDrain<{ url: string, token: string }>({
    name: 'loki',
    resolve: () => ({
      url: overrides?.url ?? process.env.LOKI_URL!,
      token: overrides?.token ?? process.env.LOKI_TOKEN!,
    }),
    encode: (events, config) => ({
      url: `${config.url}/loki/api/v1/push`,
      headers: {
        'Content-Type': 'application/json',
        Authorization: `Bearer ${config.token}`,
      },
      body: JSON.stringify({
        streams: events.map(e => ({
          stream: { service: e.service, level: e.level },
          values: [[String(Date.parse(e.timestamp) * 1e6), JSON.stringify(e)]],
        })),
      }),
    }),
  })
}

了解哪个配置优先

resolveAdapterConfig(namespace, fields, overrides) 会沿着标准链路查找,因此用户获得的配置体验与内置适配器相同:

  1. 显式传给你的工厂的 overrides
  2. runtimeConfig.evlog.<namespace>(Nuxt/Nitro)
  3. runtimeConfig.<namespace>(旧版 Nuxt/Nitro)
  4. <NS>_<FIELD> 环境变量(在 ConfigField.env 中列出 NUXT_<NS>_<FIELD>,以实现静默 Nuxt 兼容;在错误消息中仅通过 formatPublicEnvKeys 显示 <NS>_<FIELD>

字段名称应遵循项目约定:apiKeyendpointserviceNametimeout。如果你要重命名现有字段(例如 tokenapiKey),请在一个 major 版本中将两者都保留为 ConfigField 条目。有关弃用模式,请参见 axiom.tsbetter-stack.ts

将 drain 接入你的框架

一旦 createMyServiceDrain() 返回该 drain,就像其他任何东西一样接入它:

// server/plugins/evlog-drain.ts
import { createMyServiceDrain } from '~/server/utils/my-drain'

const drain = createMyServiceDrain()

export default defineNitroPlugin((nitroApp) => {
  nitroApp.hooks.hook('evlog:drain', drain)
})

对于生产环境,将其包装一次在 createDrainPipeline 中,这样事件就会被批处理并重试。

发送前进行过滤或转换

encode() 接收完整的 WideEvent[] 批次以及解析后的配置。可以内联进行过滤或转换,返回 null 是选择退出该批次的简洁方式:

encode: (events, cfg) => {
  const filtered = events.filter(e => e.level === 'error' && e.path !== '/health')
  if (filtered.length === 0) return null

  const payload = filtered.map(e => ({
    ts: new Date(e.timestamp).getTime(),
    severity: e.level.toUpperCase(),
    attributes: { method: e.method, path: e.path, status: e.status, durationMs: e.durationMs },
  }))

  return {
    url: `${cfg.endpoint}/v1/push`,
    headers: { 'Content-Type': 'application/json' },
    body: JSON.stringify(payload),
  }
}

defineDrain(非 HTTP 传输)

如果你的目标需要 gRPC、供应商 SDK、队列客户端、WebSocket 或原始套接字,可以用 defineDrain 再往下一层。传输由你负责;工具包仍然会为你提供配置解析、错误隔离和一致的形状。

import { defineDrain } from 'evlog/toolkit'

export const createCustomTransportDrain = () =>
  defineDrain<{ apiKey: string }>({
    name: 'custom',
    resolve: async () => ({ apiKey: process.env.MY_KEY! }),
    send: async (events, cfg) => {
      await myVendorSdk.publish(events, { token: cfg.apiKey })
    },
  })

send 在失败时抛出异常。defineDrain 会处理隔离:直接调用会吞掉错误并使用你的 drain 名称记录日志,而返回函数上的 raw 变体会拒绝,这正是流水线的重试和 onDropped 观察失败的方式。自行吞掉错误的 send 会使你的 drain 无法重试。

DrainContext 携带的内容

当 evlog 通过 evlog:drain 调用你的 drain 时,它会为每个事件传递一个 DrainContext

types.ts
interface DrainContext {
  /** 包含所有已累积上下文的完整宽事件 */
  event: WideEvent

  /** 请求元数据 */
  request?: {
    method: string
    path: string
    requestId: string
  }

  /** 安全的 HTTP 标头(已过滤敏感标头) */
  headers?: Record<string, string>
}

interface WideEvent {
  timestamp: string
  level: 'debug' | 'info' | 'warn' | 'error'
  service: string
  environment?: string
  version?: string
  region?: string
  commitHash?: string
  requestId?: string
  // ... 以及通过 log.set() 添加的所有字段
  [key: string]: unknown
}

encode() / send() 接收的批量形式中,你会直接得到 WideEvent[](工具包会将每个上下文中的 event 解包)。

使用工具包辅助函数

evlog/toolkit 暴露了每个内置适配器都会用到的相同辅助函数。与 drains 相关的有:

导出用途
defineHttpDrain(spec)HTTP 配方——自动重试、超时、身份标头、错误隔离
defineDrain(spec)非 HTTP 传输的相同契约
resolveAdapterConfig(ns, fields, overrides)标准配置优先级链(overrides → runtimeConfig.evlog.<ns> → env)
httpPost(opts)每个内置 HTTP 适配器使用的重试 POST 辅助函数——处理超时、重试、脱敏错误消息
composeDrains(drains)将多个 drain 合并为一个(错误隔离,并通过 Promise.allSettled 并发运行)
toTypedAttributeValue(value)将任意值转换为 Axiom / Sentry 使用的 typed attribute 形状
toOtlpAttributeValue(value)将任意值转换为 OTLP AnyValue 形状(用于 OTLP / HyperDX / PostHog 日志)
OTEL_SEVERITY_NUMBER, OTEL_SEVERITY_TEXTOTEL 日志严重级别表

向接收方标识你的流量

defineHttpDrain 会自动为每个请求添加两个头部,方便接收方识别流量:

标头
User-Agentevlog/<version>(仅限 Node / 服务器运行时——浏览器会移除此头部)
X-Evlog-Source你提供的 drain name

如果你直接基于 httpPost 构建 drain,则可以覆盖或抑制这些标头。参见身份标头

错误已经由你处理

defineHttpDrain 会自动强制执行所有最佳实践:

  1. 绝不抛出异常:捕获失败并使用 [evlog/<name>] 前缀记录日志。
  2. 重试:对于临时错误,默认尝试 2 次(可通过 retries 配置)。
  3. 超时:默认 5000ms(可通过 timeout 配置)。
  4. 优雅降级resolve() 返回 null 会使 drain 成为无操作。

如果你回退使用 defineDrain,请手动遵循相同规则。

将其发布为社区软件包

社区 drain 的推荐结构:

my-evlog-drain/
├─ src/
│  ├─ drain.ts        # 通过 defineHttpDrain 创建 createMyDrain
│  └─ index.ts        # 重新导出
├─ test/              # vitest,mock fetch
├─ package.json       # peerDependency: "evlog"
└─ README.md

evlog 添加为 peerDependency(而不是 dependency),这样你的软件包在安装时不会拉取一份 evlog 副本。

构建了很棒的东西?提交 PR,为适配器表格添加一行。社区会感谢你。

下一步