diff --git a/examples/flow/main.go b/examples/flow/main.go new file mode 100644 index 0000000..a00c8f5 --- /dev/null +++ b/examples/flow/main.go @@ -0,0 +1,26 @@ +// flow 包示例:流量统计(只统计、不限制),常配合 conn 包做连接流量上报 +package main + +import ( + "fmt" + + "git.zeroonesoft.cn/golib/zogo/flow" +) + +func main() { + f := &flow.Flow{} + + // 模拟收发流量(单位:字节) + f.Add(1024, 2048) + f.Add(512, 1024) + + send, recv := f.Get() + fmt.Printf("已发送: %d 字节, 已接收: %d 字节\n", send, recv) + + // 汇总并清零(常用于周期性上报) + send, recv = f.Reset() + fmt.Printf("Reset 返回 - 发送: %d, 接收: %d\n", send, recv) + + send, recv = f.Get() + fmt.Printf("清零后 - 发送: %d, 接收: %d\n", send, recv) +} diff --git a/flow/README.md b/flow/README.md new file mode 100644 index 0000000..ba827ea --- /dev/null +++ b/flow/README.md @@ -0,0 +1,25 @@ +# flow + +并发安全的流量统计器(只统计、不限制):按累计字节数记录收发流量, +常配合 `conn` 包给限速连接挂统计。 + +> 迁移自 go-hua/flow,代码保持原样。 + +## 用法 + +```go +import "git.zeroonesoft.cn/golib/zogo/flow" + +f := &flow.Flow{} +f.Add(100, 200) // 累计:发送 100、接收 200 +send, recv := f.Get() // (150, 280) + +send, recv = f.Reset() // 取走并清零(周期性上报场景) +``` + +完整可运行例程:[examples/flow/main.go](../examples/flow/main.go) + +## 注意 + +- 计数跨进程生命周期累计,需要"每分钟流量"这类指标时用 `Reset` 取走清零。 +- `json` 序列化字段为 `sendFlow`/`recvFlow`。 diff --git a/flow/flow.go b/flow/flow.go new file mode 100644 index 0000000..e141675 --- /dev/null +++ b/flow/flow.go @@ -0,0 +1,33 @@ +// Package flow 提供并发安全的流量统计(只统计、不限制),常配合 conn 包使用。 +package flow + +import ( + "sync/atomic" +) + +// 统计流量 不限制流量 +type Flow struct { + SendFlow int64 `json:"sendFlow"` + RecvFlow int64 `json:"recvFlow"` + FlowLimit int64 `json:"flowLimit"` +} + +// Add 累计发送/接收字节数 +func (f *Flow) Add(send, recv int64) { + atomic.AddInt64(&f.SendFlow, send) + atomic.AddInt64(&f.RecvFlow, recv) +} + +// Reset 取出并清零收发流量,返回清零前的值(用于周期性上报) +func (f *Flow) Reset() (int64, int64) { + send := atomic.SwapInt64(&f.SendFlow, 0) + recv := atomic.SwapInt64(&f.RecvFlow, 0) + return send, recv +} + +// Get 读取当前收发流量(不清零) +func (f *Flow) Get() (int64, int64) { + send := atomic.LoadInt64(&f.SendFlow) + recv := atomic.LoadInt64(&f.RecvFlow) + return send, recv +} diff --git a/flow/flow_test.go b/flow/flow_test.go new file mode 100644 index 0000000..e52b436 --- /dev/null +++ b/flow/flow_test.go @@ -0,0 +1,32 @@ +package flow + +import "testing" + +func TestFlowAddGet(t *testing.T) { + f := &Flow{} + f.Add(100, 200) + f.Add(50, 80) + + send, recv := f.Get() + if send != 150 { + t.Errorf("SendFlow = %d, want 150", send) + } + if recv != 280 { + t.Errorf("RecvFlow = %d, want 280", recv) + } +} + +func TestFlowReset(t *testing.T) { + f := &Flow{} + f.Add(10, 20) + + send, recv := f.Reset() + if send != 10 || recv != 20 { + t.Errorf("Reset 返回 = (%d, %d), want (10, 20)", send, recv) + } + + send, recv = f.Get() + if send != 0 || recv != 0 { + t.Errorf("Reset 后流量 = (%d, %d), want (0, 0)", send, recv) + } +}