README
¶
Fluent Forward Go 客户端
pkg/fluent 是一个面向 Fluentd、Fluent Bit 及其他 Forward Protocol 接收端的 Go 客户端。它使用单个后台 worker 独占连接,将编码、排队、重连、ACK 校验和关闭流程串行化,适合多个 goroutine 并发提交结构化事件。
主要能力:
- 支持 TCP、TLS/mTLS、Unix Socket、WebSocket 和安全 WebSocket;
- 支持 Message、Forward、PackedForward 和 gzip 压缩;
- 支持 Forward Protocol shared key、用户名和密码认证;
- 支持无确认投递和基于 ACK 的至少一次投递;
- 提供同步发送、有界异步队列、投递回执、队列屏障和优雅关闭;
- 所有待发送 payload 在入队前完成编码,重试期间保持内容和 chunk ID 不变。
安装
go get github.com/lwmacct/260714-go-pkg-fluent@latest
导入包:
import "github.com/lwmacct/260714-go-pkg-fluent/pkg/fluent"
快速开始
package main
import (
"context"
"log"
"time"
"github.com/lwmacct/260714-go-pkg-fluent/pkg/fluent"
)
func main() {
config := fluent.DefaultConfig("tcp://127.0.0.1:24224")
client, err := fluent.New(config)
if err != nil {
log.Fatal(err)
}
defer func() {
if err := client.Close(); err != nil {
log.Printf("关闭 Fluent 客户端失败: %v", err)
}
}()
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
err = client.Send(ctx, fluent.Event{
Tag: "app.access",
Time: time.Now(),
Record: map[string]any{
"method": "GET",
"status": 200,
},
})
if err != nil {
log.Fatal(err)
}
}
完整示例位于 examples/basic/main.go。Event.Time 为零值时,客户端会在编码时使用当前时间。
Endpoint
通过 URL 形式的 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 |
客户端默认在第一次发送时建立连接。需要在服务启动阶段提前验证连接时,可以显式调用:
if err := client.Connect(ctx); err != nil {
return err
}
连接失效后,后台 worker 会按照重试配置重新拨号;启用认证时,每次重连都会重新完成握手。
发送事件
同步发送
Send 会等待事件完成投递,或等待 context 取消:
err := client.Send(ctx, fluent.Event{
Tag: "app.audit",
Record: map[string]string{"action": "login"},
})
多个 goroutine 可以并发调用 Send,但同一客户端的网络交换由后台 worker 按入队顺序串行执行。
异步提交
Submit 等待队列出现容量,成功后返回 Receipt。Receipt.Wait 用于等待实际投递结果:
receipt, err := client.Submit(ctx, fluent.Event{
Tag: "app.event",
Record: map[string]any{"name": "created"},
})
if err != nil {
return err
}
if err := receipt.Wait(ctx); err != nil {
return err
}
传给 Submit 的 context 同时约束排队和后续投递。提交后若该 context 被取消,对应任务也会取消;Receipt.Wait 可以使用另一个 context,只限制本次等待。
非阻塞提交
TrySubmit 从不等待队列容量。队列已满时返回 fluent.ErrQueueFull:
receipt, err := client.TrySubmit(event)
switch {
case errors.Is(err, fluent.ErrQueueFull):
// 根据业务需要丢弃、计数或转存事件。
case err != nil:
return err
default:
_ = receipt
}
批量与压缩
使用 Batch 一次发送同一 Tag 下的多个事件:
err := client.SendBatch(ctx, fluent.Batch{
Tag: "app.metrics",
Encoding: fluent.BatchPackedForward,
Compression: fluent.CompressionGZIP,
Entries: []fluent.Entry{
{Time: time.Now(), Record: map[string]any{"cpu": 0.42}},
{Time: time.Now(), Record: map[string]any{"cpu": 0.57}},
},
})
编码模式:
| 值 | 行为 |
|---|---|
BatchAuto |
少于 16 条时使用 Forward,达到 16 条时使用 PackedForward |
BatchForward |
使用 Forward 模式 |
BatchPackedForward |
使用 PackedForward 模式 |
设置 CompressionGZIP 时会强制使用 gzip 压缩的 PackedForward。异步批量发送可使用 SubmitBatch,其返回值同样是 Receipt。
投递语义与 ACK
默认配置使用 DeliveryUnconfirmed:
config.Delivery = fluent.DeliveryUnconfirmed
- 不请求服务端 ACK;
- 在尚未写入任何字节时发生错误,可以安全重连并重试;
- 已经写入部分 payload 后连接失败时,不会自动重发,而是返回
ErrDeliveryUnknown; - 成功写入本地连接不代表服务端已经持久化事件。
需要更强确认时启用至少一次投递:
config.Delivery = fluent.DeliveryAtLeastOnce
- 每条消息携带随机生成的 chunk ID;
- 客户端等待并严格校验服务端 ACK;
- ACK 超时或连接中断时,使用相同 payload 和 chunk ID 重连重发;
- 如果服务端已经接收事件但 ACK 丢失,重发可能产生重复事件。
因此 DeliveryAtLeastOnce 是“至少一次”,不是 exactly-once。业务需要去重时,应在 Record 中携带业务事件 ID,并在消费端实现幂等处理。
Record 类型
Record 使用 MessagePack 编码。适合直接使用的类型包括:
map[string]any、map[string]string;- 由受支持基础类型组成的 map、slice 和 array;
- 实现
msgp.Marshaler或msgp.Encodable的类型。
普通 struct 不会自动按 json 标签反射编码。对于固定结构和高吞吐场景,建议使用 github.com/tinylib/msgp 生成编码器,让类型实现 msgp.Marshaler 或 msgp.Encodable。
Map 的键必须为字符串。Tag 不能为空,批量发送的 Entries 不能为空。
TLS 与 mTLS
TLS
config := fluent.DefaultConfig("tls://logs.example.com:24224")
config.TLSConfig = &tls.Config{
MinVersion: tls.VersionTLS13,
RootCAs: rootCAs,
}
mTLS
certificate, err := tls.LoadX509KeyPair("client.crt", "client.key")
if err != nil {
return err
}
config := fluent.DefaultConfig("tls://logs.example.com:24224")
config.TLSConfig = &tls.Config{
MinVersion: tls.VersionTLS13,
RootCAs: rootCAs,
Certificates: []tls.Certificate{certificate},
}
fluent.New 会克隆 TLS 配置;创建客户端后修改原配置不会影响已经创建的客户端。请勿在生产环境使用 InsecureSkipVerify。
Forward shared key 认证
当 Fluentd in_forward 启用安全握手时,配置 shared key:
config := fluent.DefaultConfig("tls://logs.example.com:24224")
config.Auth = &fluent.Auth{
SharedKey: []byte(os.Getenv("FLUENT_SHARED_KEY")),
Hostname: "api-01",
Username: os.Getenv("FLUENT_USERNAME"),
Password: os.Getenv("FLUENT_PASSWORD"),
}
用户名和密码只能与 shared key 一起使用。客户端会校验服务端认证结果和 PONG digest;认证拒绝可通过 errors.Is(err, fluent.ErrAuthRejected) 判断。
WebSocket
WebSocket 使用二进制消息承载 Forward payload。可设置固定 HTTP Header:
config := fluent.DefaultConfig("wss://proxy.example.com/forward")
config.WebSocket.Header = http.Header{
"X-Tenant-ID": []string{"tenant-a"},
}
需要轮换 Bearer Token 时使用 TokenProvider:
config.WebSocket.TokenProvider = func(ctx context.Context) (string, error) {
return tokenStore.Current(ctx)
}
Provider 会在每次连接尝试时调用,返回值会作为 Authorization: Bearer <token> 发送。自定义 Header 和 TLS 配置在创建客户端时都会被克隆。
生命周期
| 方法 | 语义 |
|---|---|
Connect(ctx) |
主动建立连接;不调用时由首次发送自动连接 |
Flush(ctx) |
等待调用前已入队的任务全部完成 |
Shutdown(ctx) |
停止接收新任务,排空队列并关闭连接 |
Close() |
立即取消在途操作,丢弃待处理任务并关闭连接 |
服务正常退出时优先使用 Shutdown:
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := client.Shutdown(shutdownCtx); err != nil {
_ = client.Close()
return err
}
Shutdown 或 Close 开始后,新的提交会返回 ErrClosed。Close 可以安全重复调用。
预编码消息
已有完整 Forward Protocol MessagePack payload 时,可以直接发送:
err := client.SendEncoded(ctx, fluent.EncodedMessage{Payload: payload})
客户端会校验 payload 只包含一条完整消息,并复制输入字节,调用方随后可以安全复用原切片。在 DeliveryAtLeastOnce 模式下,原始 payload 必须已经包含非空 chunk option;客户端不会修改预编码消息。
配置
建议始终从 DefaultConfig 开始修改:
config := fluent.DefaultConfig("tcp://127.0.0.1:24224")
config.TagPrefix = "production"
config.Queue.Capacity = 16384
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
默认值:
| 配置项 | 默认值 |
|---|---|
Delivery |
DeliveryUnconfirmed |
Queue.Capacity |
8192 |
Retry.MaxAttempts |
5 |
Retry.MinBackoff |
100ms |
Retry.MaxBackoff |
5s |
Timeout.Connect |
3s |
Timeout.Handshake |
3s |
Timeout.Write |
3s |
Timeout.ACK |
3s |
重试等待采用指数退避并加入抖动,最大不超过 Retry.MaxBackoff。TagPrefix 非空时,最终 Tag 为 <prefix>.<event-tag>。
高级场景可以通过 Config.Dialer 注入自定义连接:
type Dialer interface {
Dial(context.Context) (net.Conn, error)
}
设置自定义 Dialer 后,Endpoint 仅用于错误上下文,可以为空。
错误处理
可使用 errors.Is 判断稳定错误:
switch {
case errors.Is(err, fluent.ErrQueueFull):
// 异步队列已满。
case errors.Is(err, fluent.ErrClosed):
// 客户端正在或已经关闭。
case errors.Is(err, fluent.ErrDeliveryUnknown):
// 无 ACK 模式下发生部分写入,无法确定服务端是否收到。
case errors.Is(err, fluent.ErrAckMismatch):
// ACK 与当前消息的 chunk ID 不一致。
case errors.Is(err, fluent.ErrAuthRejected):
// Forward 安全握手认证失败。
}
连接或发送重试耗尽时会返回 *fluent.OpError:
var opErr *fluent.OpError
if errors.As(err, &opErr) {
log.Printf("操作=%s endpoint=%s 尝试次数=%d 临时错误=%t: %v",
opErr.Op, opErr.Endpoint, opErr.Attempt, opErr.Temporary, opErr.Err)
}
开发与验证
go test ./...
go test -race ./...
go vet ./...
golangci-lint run