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 的至少一次投递;
- 同时按未完成事件数和已编码字节数进行背压;
- 同步发送、异步回执、队列屏障、优雅关闭和立即中止;
- 结构化统计和异步失败通知;
Config承载 JSON/flag 参数,Option 注入 Go 运行时对象。
安装
go get github.com/lwmacct/260714-go-pkg-fluent@latest
import "github.com/lwmacct/260714-go-pkg-fluent/pkg/fluent"
快速开始
config := fluent.DefaultConfig()
config.Endpoint = "tcp://127.0.0.1:24224"
client, err := fluent.New(config)
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 为零值时,客户端在提交时使用当前时间;重试保留同一个时间值。
Config
Config 是所有文件友好参数的唯一模型。Endpoint、认证、队列、重试、超时、压缩、TLS 文件和静态 WebSocket Header 都直接由 Config 指定。New(config, options...) 的 Option 只用于 JSON/flag 无法表达的 Go 运行时对象。
Config.Enabled 是为应用配置和命令行 flag 预留的开关,避免调用方再包一层配置结构。pkg/fluent 不读取该字段;调用方负责根据它决定是否构造和使用客户端。
config := fluent.DefaultConfig()
config.Enabled = true
config.Endpoint = "tcp://127.0.0.1:24224"
config.TagPrefix = "production"
config.ACK = true
config.CompressionThreshold = 64 << 10
config.Buffer.MaxEvents = 16384
config.Buffer.MaxBytes = 128 << 20
config.Buffer.BatchMaxEvents = 512
config.Buffer.BatchMaxBytes = 2 << 20
config.Buffer.BatchWait = 10 * time.Millisecond
config.Retry.MaxAttempts = 8
config.Retry.MinBackoff = 200 * time.Millisecond
config.Retry.MaxBackoff = 10 * time.Second
config.Timeout.Connect = 2 * time.Second
config.Timeout.Handshake = 3 * time.Second
config.Timeout.Write = 3 * time.Second
config.Timeout.ACK = 5 * time.Second
client, err := fluent.New(config)
Config.Validate() 检查所有跨字段约束,New 会再次验证并深拷贝认证密钥、Header、TLS 配置和 Dialer。
Endpoint
| 传输方式 | Endpoint |
|---|---|
| TCP | tcp://127.0.0.1:24224 |
| TLS/mTLS | tls://logs.example.com:24224 |
| Unix Socket | unix:///run/fluent/fluent.sock |
| WebSocket | ws://proxy.example.com/forward |
| 安全 WebSocket | wss://proxy.example.com/forward |
Endpoint 不允许 user info、query 或 fragment。TCP/TLS 必须包含 host 和 port,Unix 必须使用绝对路径。与 scheme 不相容的配置会返回错误,例如 tcp:// 携带 TLS 配置、ws:// 携带 TLS 配置,或 TCP 携带 WebSocket Header。
TLS 与 mTLS
TLS 文件参数可以直接来自 JSON 或命令行 flag:
config := fluent.DefaultConfig()
config.Endpoint = "tls://logs.example.com:24224"
config.TLS.CAFile = "/etc/ssl/fluent-ca.pem"
config.TLS.CertificateFile = "/etc/ssl/client.pem"
config.TLS.KeyFile = "/etc/ssl/client-key.pem"
config.TLS.ServerName = "logs.example.com"
config.TLS.MinVersion = "1.3"
client, err := fluent.New(config)
MinVersion 支持 1.2 和 1.3。客户端证书和私钥必须同时配置。InsecureSkipVerify 只应用于显式配置,生产环境不应启用。
需要直接传入 *tls.Config 时使用 Option:
config.Endpoint = "tls://logs.example.com:24224"
tlsConfig := &tls.Config{
MinVersion: tls.VersionTLS13,
RootCAs: rootCAs,
}
client, err := fluent.New(config, fluent.WithTLSConfig(tlsConfig))
WithTLSConfig 会克隆输入配置,且不能与 Config 中的 TLS 文件字段混用。原生 *tls.Config 只通过 Option 注入,不属于文件配置模型。
WebSocket
config := fluent.DefaultConfig()
config.Endpoint = "wss://proxy.example.com/forward"
config.WebSocket.Header = http.Header{
"X-Tenant-ID": []string{"tenant-a"},
}
config.TLS.CAFile = "/etc/ssl/proxy-ca.pem"
client, err := fluent.New(
config,
fluent.WithHeaderProvider(fluent.BearerToken(tokenStore.Current)),
)
HeaderProvider 在每次连接尝试时调用,动态 Header 覆盖同名静态 Header。它是运行时函数,只通过 Option 注入。
自定义连接
代码集成使用 WithConnector 注入自定义连接器:
type Connector interface {
Dial(context.Context) (net.Conn, error)
}
config := fluent.DefaultConfig()
config.Endpoint = "service-discovery://fluent-primary" // 错误上下文标签
client, err := fluent.New(config, fluent.WithConnector(customConnector))
自定义 Connector 必须响应 context 取消。使用自定义 Connector 时,不能同时设置 TLS、WebSocket 或 NetDialer,避免两套传输配置产生歧义。
配置文件
先取得完整默认值,再将 JSON 覆盖到同一个值,未出现的字段会保留默认值:
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(config)
示例 JSON:
{
"enabled": true,
"endpoint": "tls://logs.example.com:24224",
"tag-prefix": "production",
"ack": true,
"buffer": {
"max-events": 16384,
"max-bytes": 134217728,
"batch-max-events": 512,
"batch-max-bytes": 2097152,
"batch-wait": 10000000
},
"retry": {
"max-attempts": 8,
"min-backoff": 200000000,
"max-backoff": 10000000000
},
"timeout": {
"connect": 2000000000,
"handshake": 3000000000,
"write": 3000000000,
"ack": 5000000000
},
"compression-threshold": 65536,
"tls": {
"ca-file": "/etc/ssl/fluent-ca.pem",
"server-name": "logs.example.com",
"min-version": "1.3",
"insecure-skip-verify": false
}
}
配置结构体只使用 json 和 desc tag,JSON 名称统一使用 kebab-case。desc 用作命令行 flag 的说明文本。运行时对象不进入 Config,通过 WithTLSConfig、WithConnector、WithNetDialer、WithHeaderProvider 和 WithFailureHandler 注入。Go 标准库 encoding/json 将 time.Duration 表示为纳秒整数。
发送
单条和批量共用 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,队列满时返回 ErrQueueFull。
Record
Forward Protocol 要求 Record 顶层为字符串键 map,因此公开类型固定为:
type Record map[string]any
嵌套值可以使用 MessagePack 支持的基础类型、slice、array 和字符串键 map。非法值在入队前返回 ErrInvalidRecord,不会占用队列或发起连接。
自动批处理与背压
客户端合并队首相邻、同 Tag 的异步请求,并自动选择 Message、Forward 或 PackedForward。达到 CompressionThreshold 时使用 gzip PackedForward。
Buffer.MaxEvents 和 Buffer.MaxBytes 统计全部已接收但尚未完成的工作,包括正在发送和等待 ACK 的批次。单个请求超过任一上限返回 ErrTooLarge。
投递与 ACK
默认不请求 ACK:尚未写入任何字节时可安全重试;部分写入后失败返回 ErrDeliveryUnknown,不会自动重发。
config := fluent.DefaultConfig()
config.ACK = true
client, err := fluent.New(config)
启用 ACK 后,每个 wire batch 携带随机 chunk ID。ACK 超时或连接中断会使用相同 payload 和 chunk ID 重试;格式非法或 chunk 不匹配的 ACK 不重试。ACK 丢失可能造成重复事件,业务仍需使用事件 ID 做幂等。
Forward 认证
config.Auth = &fluent.Auth{
SharedKey: 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 {
return err
}
Flush等待屏障之前接受的所有请求完成;Shutdown停止接收新任务,排空队列并关闭连接;Shutdowncontext 到期时自动调用Abort;Abort立即取消在途任务并拒绝排队任务,可重复调用;- 关闭开始后,新的提交返回
ErrClosed。
失败通知与统计
failureHandler := func(failure fluent.Failure) {
logger.Error("fluent delivery failed",
"tag", failure.Tag,
"entries", failure.Entries,
"delivered", failure.Delivered,
"error", failure.Err,
)
}
client, err := fluent.New(config, fluent.WithFailureHandler(failureHandler))
stats := client.Stats()
FailureHandler 只报告异步请求最终失败,调用在独立 goroutine 中串行执行。内部通知队列满时不阻塞投递,可通过 FailureNotificationsDropped 监控遗漏。
一个请求可能因批次上限被拆成多个 wire batch。前面的批次成功、后续失败时,Receipt 返回 *fluent.PartialDeliveryError,其中 Delivered 是已完成条目数,Total 是原请求总数。
错误处理
稳定错误可用 errors.Is 判断:ErrQueueFull、ErrTooLarge、ErrClosed、ErrInvalidRecord、ErrDeliveryUnknown、ErrACKMismatch、ErrAuthRejected、ErrServerVerification 和 ErrProtocol。
连接或发送错误带有 *fluent.Error 元数据。WebSocket HTTP upgrade 失败可通过 errors.As 取得 *fluent.WebSocketHandshakeError。
默认值
| 配置项 | 默认值 |
|---|---|
Endpoint |
tcp://127.0.0.1:24224 |
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 |
重试使用指数退避和最多 20% 的向下抖动,始终不超过 Retry.MaxBackoff。显式配置零 timeout 可禁用对应客户端超时。
开发与验证
go test ./...
go test -race ./...
go test -tags=integration ./...
go vet ./...
golangci-lint run ./...
使用 integration build tag 启用集成测试。环境中存在 fluent-bit 时,测试会启动随机本地端口验证真实 Forward TCP、ACK 和自动批处理;未安装时集成测试会 skip。