You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的死信队列实现延迟功能。

核心思路

  1. 创建两个队列:
    • delay_queue:设置10秒TTL,绑定死信交换机指向notify_queue
    • notify_queue:消费消息并发送通知
  2. 用内存Map或Redis记录userID:chatID对应的消息ID
  3. 收到message.new事件时:删除Map中对应Key的旧消息(若存在),发送新消息到delay_queue并记录消息ID
  4. 收到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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 00:29:55