Skip to content
Elogs
Esc
navigateopen⌘Jpreview
On this page

Custom Transports

把日志发到外部服务(Elasticsearch / Slack / OpenTelemetry ...)

用自定义 transport 把日志发到外部服务。transport 跟内置的 console / file 日志并行跑,统一从 emit 管道调用 —— 一条 log 命中后同时送往所有激活的 sink。

输出行为

Elogs 支持三种输出目标,互不冲突:

  1. Console 日志 —— 内置控制台输出
  2. File 日志 —— 通过 file sink 输出,详见 File Logging
  3. Transport 日志 —— 自定义外部服务,本文

开关矩阵

字段 作用
disableInternalLogger 只关 console
disableFileLogging 只关 file
useTransportsOnly 同时关 console + file,只走 transport

默认 + file(如果设置了 logFilePath)+ transports 同时生效。详见 Log Filtering

Transport 接口

Transport只有一个方法的纯接口,所有组合器、sink 出口都走它:

import type { LogLevel, Transport } from '@eastgold15/elogs'

export interface Transport {
  log: (
    level: LogLevel,                       // 'DEBUG' | 'INFO' | 'WARNING' | 'ERROR'
    message: string,
    meta?: Record<string, unknown>         // request.method/url, durationMs, pathname ...
  ) => void | Promise<void>
}

metalogToTransports 装配,固定包含 request.methodrequest.urldurationMspathname、原始 data 字段(包含 statusrequestId 等)。

基本 transport

import { Elysia } from 'elysia'
import { createElogs } from '@eastgold15/elogs'
import type { Transport } from '@eastgold15/elogs'

const consoleTransport: Transport = {
  log: (level, message, meta) => {
    console.log(`[${level}] ${message}`, meta)
  },
}

const app = new Elysia().use(
  createElogs({
    config: {
      transports: [consoleTransport],
    },
  })
)

签名是 (level, message, meta?) => void | Promise<void>。同步 transport 返回 void,异步 transport 返回 promise —— Elogs 内部会 catch 异步 reject 并节流上报(见下方错误处理)。

外部服务示例

Elasticsearch

import type { Transport } from '@eastgold15/elogs'

const elasticsearchTransport: Transport = {
  log: async (level, message, meta) => {
    await fetch('http://elasticsearch:9200/logs/_doc', {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify({
        '@timestamp': new Date().toISOString(),
        level,
        message,
        ...meta,
      }),
    })
  },
}

MongoDB

import type { Transport } from '@eastgold15/elogs'

const mongodbTransport: Transport = {
  log: async (level, message, meta) => {
    await db.collection('logs').insertOne({
      level,
      message,
      ...meta,
      timestamp: new Date(),
    })
  },
}

Slack(只告 ERROR)

import type { Transport, LogLevel } from '@eastgold15/elogs'

const slackTransport: Transport = {
  log: async (level: LogLevel, message, meta) => {
    if (level !== 'ERROR') return
    const webhook = process.env.SLACK_WEBHOOK_URL
    if (!webhook) return

    await fetch(webhook, {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify({
        text: `[${level}] ${message}\n\`\`\`${JSON.stringify(meta, null, 2)}\`\`\``,
      }),
    })
  },
}

meta 已经是纯 Record<string, unknown>,可以直接 JSON.stringify 序列化。

OpenTelemetry / Datadog / Loki

任何 HTTP 端点都能直接 fetch 包一个 transport。批量 / 节流需求用下面的 batch 组合器:

import type { Transport } from '@eastgold15/elogs'

const otelTransport: Transport = {
  log: async (level, message, meta) => {
    await fetch('https://otel-collector:4318/v1/logs', {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify({
        severityText: level,
        body: message,
        attributes: meta,
        time: new Date().toISOString(),
      }),
    })
  },
}

多个 transport

app.use(
  createElogs({
    config: {
      transports: [elasticsearchTransport, slackTransport],
    },
  })
)

数组里的 transport 顺序调用,各自独立 —— Elasticsearch 慢不会阻塞 Slack。每个 transport 的错误也独立节流,一个挂了不影响其它 transport 的错误日志。

只用 transport

app.use(
  createElogs({
    config: {
      useTransportsOnly: true,           // 关掉 console + file
      transports: [elasticsearchTransport, slackTransport],
    },
  })
)

生产只接外部聚合器时的标准姿势 —— 避免本地双写,也防止容器没挂 volume 时日志全丢。

错误处理 & 节流

transport 抛出的同步异常和异步 reject 都会被 logToTransports 捕获并打印到 console.error,不会让请求挂掉。但为了防止坏掉的 transport 把 stderr 刷爆,同一 transport 的同一错误在 transportThrottleMs(默认 5000ms)窗口内合并为一条 console.error

createElogs({
  config: {
    transportThrottleMs: 10_000, // 10s 窗口
  },
})

这是生产环境的安全阀 —— 你的 transport 仍要自己处理错误(重试 / 降级 / 死信队列),不要依赖 Elogs 的内部日志。

最佳实践

  • 在 transport 内部处理错误 —— Elogs 捕获了 throw 但只打 console.error,不会送回外部目标。重试、熔断、死信队列都在 transport 里实现。
  • 异步用 Promise<void> —— 同步 transport 用 void,异步返回 Promise<void> 让 Elogs 知道要 .catch
  • 避免在 meta 里写 secret / PII —— autoRedact 处理 message 字符串和 key 名,但 transport 拿到的 meta 是 emit 管道给的原始数据,不会二次脱敏。敏感字段在源头就过滤掉。
  • 生产只接外部聚合器时,设 useTransportsOnly: true —— 避免双写也避免本地文件膨胀。
  • 按重要性排 transport 顺序 —— 告警(快速失败)放最前,批写(慢)放最后。
  • 慢 transport 套 batch —— 远程 HTTP 每条一次请求成本压垮服务,见下。

组合器(Composers)

Transport 接口只有一个 log 方法,所以 (Transport) => Transport 的纯函数可以任意堆叠,跟 pino multi-stream 一样。Elogs 给 5 个常用组合器,从 @eastgold15/elogs 直接 import:

组合器 用途
tee([t1, t2, ...]) 多路 fan-out,任一失败不短路
sample(rate, t) 概率采样(0-1)
filter(predicate, t) 谓词过滤
tap(fn, t) 旁路(metrics / 计数),不影响主链路
batch(size, ms, t) 缓冲 size 条或 ms 毫秒,合成一条发

所有组合器都在 output/composers.ts 里实现,签名见 LogEntry

tee —— 多路 fan-out

把一条 log 同时发到多个 transport。任一分支 throw 或 reject 都不会短路其它分支(Promise.allSettled 兜底)。

import { createElogs, tee } from '@eastgold15/elogs'
import type { Transport } from '@eastgold15/elogs'

const consoleT: Transport = { log: (l, m) => console.log(l, m) }
const metricsT: Transport = { log: (l, m, meta) => metrics.increment(m) }

createElogs({
  config: {
    transports: [tee([consoleT, metricsT])],
  },
})

sample —— 概率采样

rate (0-1) 的概率把 log 转发给 transport。rate=0 完全丢弃,rate=1 等于直通。rate 越界抛 Error

import { createElogs, tee, sample } from '@eastgold15/elogs'

createElogs({
  config: {
    transports: [
      tee([
        sample(0.1, metricsT),  // 10% 进 metrics
        consoleT,                // 100% 进 console
      ]),
    ],
  },
})

filter —— 谓词过滤

仅当 predicate(level, message, meta) 返回 truthy 时才转发。最常见 才告警。

import { createElogs, tee, filter, sample } from '@eastgold15/elogs'
import type { Transport, LogLevel } from '@eastgold15/elogs'

const slackAlerter: Transport = {
  log: async (l: LogLevel, m) => {
    await fetch(process.env.SLACK_WEBHOOK!, {
      method: 'POST',
      body: JSON.stringify({ text: `[${l}] ${m}` }),
    })
  },
}

createElogs({
  config: {
    transports: [
      tee([
        filter((lvl) => lvl === 'ERROR', slackAlerter),
        sample(0.1, metricsT),
        consoleT,
      ]),
    ],
  },
})

tap —— 旁路

对每条 log 跑 fn,然后总是转发给 transport。适合打 metrics / 计数 / 调试 trace。fn 抛错会被静默吞(不影响主链路)。

import { createElogs, tap, sample } from '@eastgold15/elogs'

let counter = 0

createElogs({
  config: {
    transports: [
      tap(
        (lvl, msg, meta) => {
          counter++
          if (meta?.durationMs && meta.durationMs > 1000) {
            console.warn('slow request', meta)
          }
        },
        sample(0.01, metricsT),  // 计数 100%,metrics 采样 1%
      ),
    ],
  },
})

batch —— 缓冲

缓冲 log 直到满 size距首条 flushMs 毫秒,触发后作为一条 INFO “batch” 发出去,meta.entries 是缓冲数组。size / flushMs ≤ 0 抛 Errorbatch 内部会暴露 flush 方法(在返回的 transport 上 cast 出来),process.on("beforeExit", ...) / SIGTERM handler 里调一下能避免退出时丢 buffer。

适合打远程服务(每次 HTTP 请求成本高):

import { createElogs, batch, filter, tee } from '@eastgold15/elogs'
import type { Transport, LogLevel } from '@eastgold15/elogs'

const httpSink: Transport = {
  log: async (lvl: LogLevel, msg, meta) => {
    await fetch('https://logs.example.com/ingest', {
      method: 'POST',
      body: JSON.stringify({ lvl, msg, meta }),
    })
  },
}

const debugBatch = batch(50, 5000, httpSink)

process.on('beforeExit', () => {
  ;(debugBatch as Transport & { flush?: () => void }).flush?.()
})

createElogs({
  config: {
    transports: [
      tee([
        filter((lvl) => lvl === 'DEBUG', debugBatch),  // DEBUG 批量 5s
        consoleT,
      ]),
    ],
  },
})

组合原则

  • 每个组合器只做一件事,任意堆叠;
  • 顺序敏感:filter 放在 tee 之前比之后便宜(早 fail 早 return);
  • 涉及副作用的(metrics / 计数)用 tap,绝不要filter 丢日志;
  • 远程 / 慢 transport 一定套 batch,否则一个请求一条 HTTP 成本压垮服务;
  • batchbeforeExit / SIGTERM 之前手动 flush,否则最后一批进 buffer 的会丢。

API 参考

Last updated on August 15, 2026

Was this page helpful?