260714-go-pkg-fluent

module
v0.3.260716 Latest Latest
Warning

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

Go to latest
Published: Jul 15, 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 的至少一次投递;
  • 同时按未完成事件数和已编码字节数进行背压;
  • 同步发送、异步回执、队列屏障、优雅关闭和立即中止;
  • 结构化统计和异步失败通知。

安装

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.goEntry.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),
)

MaxEventsMaxBytes 统计全部已接收但尚未完成的工作,包括正在发送和等待 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 停止接收新任务,排空队列并关闭连接;
  • Shutdown context 到期时会调用 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 映射的配置字段:AuthTagPrefixACKBufferRetryTimeoutCompressionThreshold。推荐先取得完整默认值,再将配置文件覆盖到同一个值:

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),
)

配置结构体只使用 jsondesc tag;desc 用作命令行 flag 的说明文本。FailureHandler 是运行时函数,使用 json:"-" desc:"-" 跳过配置文件和 flag。Go 标准库 encoding/jsontime.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 提供当前 PendingEventsPendingBytes,以及累计提交、成功、失败、重试、连接和拒绝计数。

错误处理

稳定错误可用 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,读取状态码和受限响应正文。

一个较大的请求可能因 BatchMaxEventsBatchMaxBytes 被拆成多个 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。

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