260714-go-pkg-fluent

module
v0.6.260716 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 16, 2026 License: MIT

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.goEntry.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.21.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 时,不能同时设置 TLSWebSocketNetDialer,避免两套传输配置产生歧义。

配置文件

先取得完整默认值,再将 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
  }
}

配置结构体只使用 jsondesc tag,JSON 名称统一使用 kebab-case。desc 用作命令行 flag 的说明文本。运行时对象不进入 Config,通过 WithTLSConfigWithConnectorWithNetDialerWithHeaderProviderWithFailureHandler 注入。Go 标准库 encoding/jsontime.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.MaxEventsBuffer.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 停止接收新任务,排空队列并关闭连接;
  • Shutdown context 到期时自动调用 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 判断:ErrQueueFullErrTooLargeErrClosedErrInvalidRecordErrDeliveryUnknownErrACKMismatchErrAuthRejectedErrServerVerificationErrProtocol

连接或发送错误带有 *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。

Directories

Path Synopsis
examples
basic command
pkg
fluent
Package fluent sends structured events using Fluent Forward Protocol v1.
Package fluent sends structured events using Fluent Forward Protocol v1.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL