| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869 |
- package lc
- import (
- "sync"
- "time"
- )
- // CommandQueue 命令队列,用于串行处理Send操作(带150ms间隔)
- type CommandQueue struct {
- tasks chan []byte // 存储待发送的数据
- running bool // 队列运行状态
- mu sync.Mutex // 保护running状态的互斥锁
- wg sync.WaitGroup // 用于等待队列结束
- interval time.Duration // 任务执行间隔(毫秒)
- }
- // NewCommandQueue 创建新的命令队列(指定间隔,默认150ms)
- func NewCommandQueue(bufferSize int, interval ...time.Duration) *CommandQueue {
- // 默认间隔150毫秒,支持自定义传入
- intervalMs := 200 * time.Millisecond
- if len(interval) > 0 {
- intervalMs = interval[0]
- }
- return &CommandQueue{
- tasks: make(chan []byte, bufferSize),
- interval: intervalMs,
- }
- }
- // Start 启动队列处理循环
- func (q *CommandQueue) Start(handle func([]byte)) {
- q.mu.Lock()
- defer q.mu.Unlock()
- if q.running {
- return
- }
- q.running = true
- q.wg.Add(1)
- go func() {
- defer q.wg.Done()
- // 循环处理队列中的发送任务
- for data := range q.tasks {
- handle(data) // 执行实际发送操作
- time.Sleep(q.interval)
- }
- }()
- }
- // AddTask 添加发送任务到队列
- func (q *CommandQueue) AddTask(data []byte) {
- q.mu.Lock()
- defer q.mu.Unlock()
- if !q.running {
- return
- }
- q.tasks <- data
- }
- // Stop 停止队列,等待所有任务处理完毕
- func (q *CommandQueue) Stop() {
- q.mu.Lock()
- defer q.mu.Unlock()
- if !q.running {
- return
- }
- close(q.tasks) // 关闭通道,让处理goroutine退出
- q.wg.Wait() // 等待所有任务处理完毕
- q.running = false
- }
|