README
¶
Fluent Forward Go 客户端
pkg/fluent 是一个面向 Fluentd、Fluent Bit 及其他 Forward Protocol v1 接收端的 Go 客户端。客户端使用有界内存队列、自动批处理和单连接 worker,适合多个 goroutine 并发提交结构化事件。
主要能力:
- TCP、TLS/mTLS、Unix Socket、WebSocket 和安全 WebSocket;
- 自动选择 Message、Forward、PackedForward 和 gzip;
- Forward shared key、用户名和密码认证;
- 无确认投递和基于 ACK 的至少一次投递;
- 同时按未完成事件数和已编码字节数进行背压;
- 同步发送、异步回执、队列屏障、优雅关闭和立即中止;
- 结构化统计和异步失败通知。
安装
go get github.com/lwmacct/260714-go-pkg-fluent@latest
import "github.com/lwmacct/260714-go-pkg-fluent/pkg/fluent"
快速开始
client, err := fluent.New(fluent.TCPConnector{
Address: "127.0.0.1:24224",
})
if err != nil {
return err
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
err = client.Send(ctx, "app.access", fluent.Entry{
Time: time.Now(),
Record: fluent.Record{
"method": "GET",
"status": 200,
},
})
if err != nil {
client.Abort()
return err
}
if err := client.Shutdown(ctx); err != nil {
return err
}
完整示例位于 examples/basic/main.go。Entry.Time 为零值时,客户端会在提交时使用当前时间;重试会保留同一个时间值。
Connector
传输通过互斥的 Connector 类型配置,不存在与当前传输无关、被静默忽略的字段。
TCP
client, err := fluent.New(fluent.TCPConnector{
Address: "127.0.0.1:24224",
})
TLS 与 mTLS
client, err := fluent.New(fluent.TCPConnector{
Address: "logs.example.com:24224",
TLS: &tls.Config{
MinVersion: tls.VersionTLS13,
RootCAs: rootCAs,
Certificates: []tls.Certificate{certificate},
},
})
Unix Socket
client, err := fluent.New(fluent.UnixConnector{
Path: "/run/fluent/fluent.sock",
})
WebSocket
client, err := fluent.New(fluent.WebSocketConnector{
URL: "wss://proxy.example.com/forward",
Header: http.Header{
"X-Tenant-ID": []string{"tenant-a"},
},
HeaderProvider: fluent.BearerToken(tokenStore.Current),
TLS: tlsConfig,
})
HeaderProvider 在每次连接尝试时调用,动态 Header 覆盖同名静态 Header。内建 Connector 的 TLS 配置、HTTP Header 和 net.Dialer 会在创建客户端时复制。
也可以实现自定义 Connector:
type Connector interface {
Dial(context.Context) (net.Conn, error)
}
实现必须让完整建连过程响应 context 取消。TimeoutConfig.Connect 统一包裹内建和自定义 Connector。
发送
单条和批量共用 Send:
err := client.Send(ctx, "app.metrics",
fluent.Entry{Record: fluent.Record{"cpu": 0.42}},
fluent.Entry{Record: fluent.Record{"cpu": 0.57}},
)
Send 的 context 同时控制等待队列容量、实际投递和结果等待。同一个 Client 的网络交换串行执行,多个 goroutine 可以安全并发调用。
异步提交
receipt, err := client.Submit(ctx, "app.audit", fluent.Entry{
Record: fluent.Record{"action": "login"},
})
if err != nil {
return err
}
if err := receipt.Wait(waitCtx); err != nil {
return err
}
Submit 的 context 只控制等待队列容量。一旦成功返回,客户端接管事件,调用方 context 取消不会撤销投递。Receipt.Wait 的 context 也只控制本次等待。
非阻塞提交使用 TrySubmit:
receipt, err := client.TrySubmit("app.event", entry)
switch {
case errors.Is(err, fluent.ErrQueueFull):
// 丢弃、计数或转存。
case err != nil:
return err
default:
_ = receipt
}
不要无条件丢弃 Receipt。需要 fire-and-forget 时应配置 WithFailureHandler,保证最终失败可见。
Record
Forward Protocol 要求 Record 顶层为字符串键 map,因此公开类型固定为:
type Record map[string]any
嵌套值可以使用 MessagePack 支持的基础类型、slice、array 和字符串键 map。非法值会在入队前返回 ErrInvalidRecord,不会占用队列或发起连接。普通 struct 不会自动按 JSON 标签反射编码。
自动批处理与压缩
客户端会合并队首相邻、同 Tag 的异步请求,并自动选择编码方式:
- 单条使用 Message;
- 小批量使用 Forward;
- 大批量或较大 payload 使用 PackedForward;
- 达到配置阈值时使用 gzip PackedForward。
调用方无需理解或选择 wire encoding。同步 Send 不会与其他调用方合并,确保其 context 取消只影响自己的请求。
client, err := fluent.New(
connector,
fluent.WithBuffer(fluent.BufferConfig{
MaxEvents: 8192,
MaxBytes: 64 << 20,
BatchMaxEvents: 256,
BatchMaxBytes: 1 << 20,
BatchWait: 5 * time.Millisecond,
}),
fluent.WithCompression(64<<10),
)
MaxEvents 和 MaxBytes 统计全部已接收但尚未完成的工作,包括正在发送和等待 ACK 的批次。单个请求超过任一上限会返回 ErrTooLarge。
投递语义与 ACK
默认不请求 ACK:
- 尚未写入任何字节时失败,可安全重连重试;
- 部分写入后失败返回
ErrDeliveryUnknown,不会自动重发; - 写入本地连接成功不代表服务端已经持久化。
启用至少一次投递:
client, err := fluent.New(connector, fluent.WithACK())
- 每个实际批次携带随机 chunk ID;
- ACK 超时或连接中断时,使用完全相同的 payload 和 chunk ID 重试;
- 格式非法或 chunk 不匹配的 ACK 不会重试;
- ACK 丢失可能造成重复事件,因此业务仍需使用事件 ID 做幂等。
Forward 认证
client, err := fluent.New(
connector,
fluent.WithAuth(fluent.Auth{
SharedKey: []byte(os.Getenv("FLUENT_SHARED_KEY")),
Hostname: "api-01",
Username: os.Getenv("FLUENT_USERNAME"),
Password: os.Getenv("FLUENT_PASSWORD"),
}),
)
用户名和密码必须与 shared key 一起使用。认证拒绝返回 ErrAuthRejected;服务端 digest 校验失败返回 ErrServerVerification。
生命周期
if err := client.Flush(ctx); err != nil {
return err
}
if err := client.Shutdown(ctx); err != nil {
// Shutdown 超时已自动中止剩余任务。
return err
}
Flush等待屏障之前接受的所有请求完成;Shutdown停止接收新任务,排空队列并关闭连接;Shutdowncontext 到期时会调用Abort,所以返回时间有明确上界;Abort立即取消在途任务并拒绝排队任务,可安全重复调用;- 生命周期关闭开始后,新的提交返回
ErrClosed。
超时与重试
client, err := fluent.New(
connector,
fluent.WithTimeouts(fluent.TimeoutConfig{
Connect: 2 * time.Second,
Handshake: 3 * time.Second,
Write: 3 * time.Second,
ACK: 5 * time.Second,
}),
fluent.WithRetry(fluent.RetryConfig{
MaxAttempts: 8,
MinBackoff: 200 * time.Millisecond,
MaxBackoff: 10 * time.Second,
}),
)
显式传入 TimeoutConfig 时,单项零值表示禁用该超时。重试采用指数退避和最多 20% 的向下抖动,始终不超过 MaxBackoff。
配置文件
Config 包含可由 JSON 映射的配置字段:Auth、TagPrefix、ACK、Buffer、Retry、Timeout 和 CompressionThreshold。推荐先取得完整默认值,再将配置文件覆盖到同一个值:
config := fluent.DefaultConfig()
if err := json.Unmarshal(data, &config); err != nil {
return err
}
if err := config.Validate(); err != nil {
return err
}
client, err := fluent.New(connector, config)
这样配置文件中未出现的字段会保留默认值。Config 本身实现 Option,后续 Option 可以继续覆盖运行时字段:
client, err := fluent.New(
connector,
config,
fluent.WithFailureHandler(reportFailure),
)
配置结构体只使用 json 和 desc tag;desc 用作命令行 flag 的说明文本。FailureHandler 是运行时函数,使用 json:"-" desc:"-" 跳过配置文件和 flag。Go 标准库 encoding/json 将 time.Duration 表示为纳秒整数。
默认值:
| 配置项 | 默认值 |
|---|---|
Buffer.MaxEvents |
8192 |
Buffer.MaxBytes |
64MiB |
Buffer.BatchMaxEvents |
256 |
Buffer.BatchMaxBytes |
1MiB |
Buffer.BatchWait |
5ms |
Retry.MaxAttempts |
5 |
Retry.MinBackoff |
100ms |
Retry.MaxBackoff |
5s |
Timeout.Connect |
3s |
Timeout.Handshake |
3s |
Timeout.Write |
3s |
Timeout.ACK |
3s |
失败通知与统计
client, err := fluent.New(
connector,
fluent.WithFailureHandler(func(failure fluent.Failure) {
logger.Error("fluent delivery failed",
"tag", failure.Tag,
"entries", failure.Entries,
"delivered", failure.Delivered,
"error", failure.Err,
)
}),
)
stats := client.Stats()
FailureHandler 只报告异步请求最终失败,调用在独立 goroutine 中串行执行。内部通知队列满时不会阻塞投递,可通过 FailureNotificationsDropped 监控遗漏。
Stats 提供当前 PendingEvents、PendingBytes,以及累计提交、成功、失败、重试、连接和拒绝计数。
错误处理
稳定错误可用 errors.Is 判断:
switch {
case errors.Is(err, fluent.ErrQueueFull):
case errors.Is(err, fluent.ErrTooLarge):
case errors.Is(err, fluent.ErrClosed):
case errors.Is(err, fluent.ErrInvalidRecord):
case errors.Is(err, fluent.ErrDeliveryUnknown):
case errors.Is(err, fluent.ErrACKMismatch):
case errors.Is(err, fluent.ErrAuthRejected):
case errors.Is(err, fluent.ErrServerVerification):
case errors.Is(err, fluent.ErrProtocol):
}
连接或发送错误带有 *fluent.Error 元数据:
var operationErr *fluent.Error
if errors.As(err, &operationErr) {
log.Printf("operation=%s endpoint=%s attempt=%d retryable=%t: %v",
operationErr.Operation,
operationErr.Endpoint,
operationErr.Attempt,
operationErr.Retryable,
operationErr.Err,
)
}
WebSocket HTTP upgrade 失败可通过 errors.As 取得 *fluent.WebSocketHandshakeError,读取状态码和受限响应正文。
一个较大的请求可能因 BatchMaxEvents 或 BatchMaxBytes 被拆成多个 wire batch。前面的批次成功、后续批次失败时,Receipt 返回 *fluent.PartialDeliveryError,其中 Delivered 是已经完成的条目数,Total 是原请求总数;FailureHandler 的 Delivered 字段提供相同信息。
开发与验证
go test ./...
go test -race ./...
go vet ./...
golangci-lint run ./...
环境中可以找到 fluent-bit 时,go test ./... 还会启动一个随机本地端口,验证真实 Forward TCP、ACK 和自动批处理;未安装时该集成测试会 skip。