command_queue.go 1.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869
  1. package lc
  2. import (
  3. "sync"
  4. "time"
  5. )
  6. // CommandQueue 命令队列,用于串行处理Send操作(带150ms间隔)
  7. type CommandQueue struct {
  8. tasks chan []byte // 存储待发送的数据
  9. running bool // 队列运行状态
  10. mu sync.Mutex // 保护running状态的互斥锁
  11. wg sync.WaitGroup // 用于等待队列结束
  12. interval time.Duration // 任务执行间隔(毫秒)
  13. }
  14. // NewCommandQueue 创建新的命令队列(指定间隔,默认150ms)
  15. func NewCommandQueue(bufferSize int, interval ...time.Duration) *CommandQueue {
  16. // 默认间隔150毫秒,支持自定义传入
  17. intervalMs := 200 * time.Millisecond
  18. if len(interval) > 0 {
  19. intervalMs = interval[0]
  20. }
  21. return &CommandQueue{
  22. tasks: make(chan []byte, bufferSize),
  23. interval: intervalMs,
  24. }
  25. }
  26. // Start 启动队列处理循环
  27. func (q *CommandQueue) Start(handle func([]byte)) {
  28. q.mu.Lock()
  29. defer q.mu.Unlock()
  30. if q.running {
  31. return
  32. }
  33. q.running = true
  34. q.wg.Add(1)
  35. go func() {
  36. defer q.wg.Done()
  37. // 循环处理队列中的发送任务
  38. for data := range q.tasks {
  39. handle(data) // 执行实际发送操作
  40. time.Sleep(q.interval)
  41. }
  42. }()
  43. }
  44. // AddTask 添加发送任务到队列
  45. func (q *CommandQueue) AddTask(data []byte) {
  46. q.mu.Lock()
  47. defer q.mu.Unlock()
  48. if !q.running {
  49. return
  50. }
  51. q.tasks <- data
  52. }
  53. // Stop 停止队列,等待所有任务处理完毕
  54. func (q *CommandQueue) Stop() {
  55. q.mu.Lock()
  56. defer q.mu.Unlock()
  57. if !q.running {
  58. return
  59. }
  60. close(q.tasks) // 关闭通道,让处理goroutine退出
  61. q.wg.Wait() // 等待所有任务处理完毕
  62. q.running = false
  63. }