260714-go-pkg-fluent

module
v0.2.260714 Latest Latest
Warning

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

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

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.goEvent.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 等待队列出现容量,成功后返回 ReceiptReceipt.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]anymap[string]string
  • 由受支持基础类型组成的 map、slice 和 array;
  • 实现 msgp.Marshalermsgp.Encodable 的类型。

普通 struct 不会自动按 json 标签反射编码。对于固定结构和高吞吐场景,建议使用 github.com/tinylib/msgp 生成编码器,让类型实现 msgp.Marshalermsgp.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
}

ShutdownClose 开始后,新的提交会返回 ErrClosedClose 可以安全重复调用。

预编码消息

已有完整 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.MaxBackoffTagPrefix 非空时,最终 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

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