Golang实现延迟10秒推送通知任务及事件管理方案咨询
基于Golang的延迟推送通知解决方案
针对你的需求——处理message.new和message.read事件、支持同一用户/聊天的新消息替换旧事件、10秒延迟后发送通知,以下是几个轻量且贴合Golang生态的实现方案:
方案1:时间轮(Timing Wheel)内存实现
这是最匹配需求的方案,既能高效处理延迟任务,又能快速完成任务的替换和删除操作,且无外部依赖。
核心思路
- 用一个内存Map存储待执行任务的映射,Key设为
userID:chatID(唯一标识同一用户的同一聊天会话),Value为时间轮中的任务句柄 - 收到
message.new事件时:检查Map中是否存在对应Key的任务,存在则取消旧任务并删除映射;然后添加新任务到时间轮,10秒后触发通知发送,执行完成后自动移除Map中的记录 - 收到
message.read事件时:根据userID:chatID查找Map中的任务,存在则取消任务并删除映射
代码示例(基于RussellLuo/timingwheel库)
package main import ( "fmt" "time" "github.com/RussellLuo/timingwheel" ) // 定义事件结构 type MessageEvent struct { UserID string ChatID string Content string } // 模拟通知服务 type NotificationService struct{} func (ns *NotificationService) Push(event MessageEvent) error { fmt.Printf("[通知发送] 用户%s 聊天%s:%s\n", event.UserID, event.ChatID, event.Content) return nil } func main() { // 初始化时间轮:刻度100ms,总刻度数100(覆盖10秒延迟需求) tw := timingwheel.NewTimingWheel(100*time.Millisecond, 100) tw.Start() defer tw.Stop() // 存储待执行任务的映射 pendingTasks := make(map[string]*timingwheel.Task) notifyService := &NotificationService{} // 处理message.new事件 handleNewMessage := func(event MessageEvent) { key := fmt.Sprintf("%s:%s", event.UserID, event.ChatID) // 替换旧任务 if oldTask, exists := pendingTasks[key]; exists { oldTask.Cancel() delete(pendingTasks, key) } // 添加新的延迟任务 newTask := tw.AfterFunc(10*time.Second, func() { notifyService.Push(event) delete(pendingTasks, key) }) pendingTasks[key] = newTask } // 处理message.read事件 handleReadMessage := func(userID, chatID string) { key := fmt.Sprintf("%s:%s", userID, chatID) if task, exists := pendingTasks[key]; exists { task.Cancel() delete(pendingTasks, key) fmt.Printf("[任务取消] 用户%s 聊天%s的通知已取消\n", userID, chatID) } } // 模拟业务场景 // 1. 用户1收到第一条消息 handleNewMessage(MessageEvent{UserID: "user001", ChatID: "chat001", Content: "你好,第一条消息"}) // 2. 5秒后同一用户收到第二条消息,替换旧任务 time.Sleep(5 * time.Second) handleNewMessage(MessageEvent{UserID: "user001", ChatID: "chat001", Content: "你好,更新后的消息"}) // 3. 3秒后用户读取消息,取消通知 time.Sleep(3 * time.Second) handleReadMessage("user001", "chat001") // 等待程序结束 time.Sleep(10 * time.Second) }
优缺点
- ✅ 纯内存实现,轻量高效,延迟精度高
- ✅ 支持快速的任务查找、替换和删除,完全匹配需求
- ❌ 服务重启会丢失未执行任务,若需持久化可结合LevelDB等轻量KV存储做状态备份
方案2:RabbitMQ延迟队列(带持久化)
如果需要任务持久化、避免重启丢失,可选择RabbitMQ的死信队列实现延迟功能。
核心思路
- 创建两个队列:
delay_queue:设置10秒TTL,绑定死信交换机指向notify_queuenotify_queue:消费消息并发送通知
- 用内存Map或Redis记录
userID:chatID对应的消息ID - 收到
message.new事件时:删除Map中对应Key的旧消息(若存在),发送新消息到delay_queue并记录消息ID - 收到
message.read事件时:根据Key取出消息ID,删除RabbitMQ中的对应消息并移除Map记录
优缺点
- ✅ 自带持久化,任务可靠性高
- ✅ 成熟的消息队列生态,便于扩展
- ❌ 需要部署维护RabbitMQ,增加外部依赖
- ❌ 消息删除操作相对繁琐,需跟踪消息ID
方案3:轻量调度库(gocron)实现
用go-co-op/gocron这类轻量调度库,结合内存Map快速实现需求。
代码示例
package main import ( "fmt" "time" "github.com/go-co-op/gocron" ) type MessageEvent struct { UserID string ChatID string Content string } type NotificationService struct{} func (ns *NotificationService) Push(event MessageEvent) error { fmt.Printf("[通知发送] 用户%s 聊天%s:%s\n", event.UserID, event.ChatID, event.Content) return nil } func main() { scheduler := gocron.NewScheduler(time.UTC) scheduler.StartAsync() defer scheduler.Stop() pendingJobs := make(map[string]*gocron.Job) notifyService := &NotificationService{} handleNewMessage := func(event MessageEvent) { key := fmt.Sprintf("%s:%s", event.UserID, event.ChatID) if oldJob, exists := pendingJobs[key]; exists { oldJob.Cancel() delete(pendingJobs, key) } job, _ := scheduler.Every(10).Seconds().LimitRunsTo(1).Do(func() { notifyService.Push(event) delete(pendingJobs, key) }) pendingJobs[key] = job } handleReadMessage := func(userID, chatID string) { key := fmt.Sprintf("%s:%s", userID, chatID) if job, exists := pendingJobs[key]; exists { job.Cancel() delete(pendingJobs, key) fmt.Printf("[任务取消] 用户%s 聊天%s的通知已取消\n", userID, chatID) } } // 模拟场景 handleNewMessage(MessageEvent{UserID: "user002", ChatID: "chat002", Content: "测试消息"}) time.Sleep(6 * time.Second) handleNewMessage(MessageEvent{UserID: "user002", ChatID: "chat002", Content: "更新测试消息"}) time.Sleep(2 * time.Second) handleReadMessage("user002", "chat002") time.Sleep(10 * time.Second) }
优缺点
- ✅ API友好,实现成本低
- ✅ 支持任务取消和替换
- ❌ 纯内存存储,重启丢失任务,需额外做持久化
方案选择建议
- 若无需持久化、追求极致轻量高效:优先选择时间轮方案
- 若需任务持久化、可靠性要求高:选择RabbitMQ延迟队列
- 若快速原型开发、对持久化要求不高:可选用gocron调度库方案
内容的提问来源于stack exchange,提问作者alex
相关产品推荐
相关产品推荐

