Go Fiber预订API推送通知方案:解决重启丢失定时任务问题
解决方案:可靠定时推送与持久化任务调度
针对你的Go Fiber预订API定时推送需求,核心问题是解决服务器重启后定时任务丢失的问题,以下是可落地、可扩展的方案:
一、核心思路:放弃内存定时,采用持久化任务调度
time.AfterFunc()依赖内存存储任务,服务器重启必然丢失。可靠的方案需要将待执行的推送任务持久化到数据库,通过后台调度器定期扫描并执行,重启后自动恢复未完成的任务。
1. 扩展数据模型
新增NotificationTask集合(表)存储推送任务,关联预订ID并记录关键信息:
type NotificationTask struct { ID primitive.ObjectID `bson:"_id,omitempty"` BookingID primitive.ObjectID `bson:"booking_id"` Type string `bson:"type"` // 可选值:remind_10min/started/ended TriggerAt time.Time `bson:"trigger_at"` // 推送触发时间(UTC) Status string `bson:"status"` // 可选值:pending/executed/failed/canceled RetryCount int `bson:"retry_count"`// 失败重试次数 }
2. 任务生成逻辑
在创建预订时,直接计算三类推送的触发时间并插入任务:
func createBooking(c *fiber.Ctx) error { var booking Booking if err := c.BodyParser(&booking); err != nil { return err } // 初始化预订基础字段 booking.ID = primitive.NewObjectID() booking.Created = time.Now().UTC() booking.Active = false // 保存预订到数据库 if _, err := bookingsCollection.InsertOne(c.Context(), booking); err != nil { return err } // 生成三个推送任务 tasks := []NotificationTask{ { BookingID: booking.ID, Type: "remind_10min", TriggerAt: booking.StartDate.UTC().Add(-10 * time.Minute), Status: "pending", }, { BookingID: booking.ID, Type: "started", TriggerAt: booking.StartDate.UTC(), Status: "pending", }, { BookingID: booking.ID, Type: "ended", TriggerAt: booking.EndDate.UTC(), Status: "pending", }, } if _, err := notificationTasksCollection.InsertMany(c.Context(), tasks); err != nil { return err } return c.JSON(booking) }
3. 后台调度器实现
启动一个常驻goroutine,定期扫描待执行任务并触发推送:
func startTaskScheduler(ctx context.Context) { // 每分钟扫描一次,可根据业务调整频率 ticker := time.NewTicker(1 * time.Minute) defer ticker.Stop() for { select { case <-ticker.C: now := time.Now().UTC() // 查询当前时间之前的待执行任务 cursor, err := notificationTasksCollection.Find(ctx, bson.M{ "trigger_at": bson.M{"$lte": now}, "status": "pending", }) if err != nil { log.Printf("获取待执行任务失败: %v", err) continue } var tasks []NotificationTask if err := cursor.All(ctx, &tasks); err != nil { log.Printf("解析任务失败: %v", err) continue } // 逐个执行推送 for _, task := range tasks { err := sendExpoNotification(task.BookingID, task.Type) if err != nil { log.Printf("推送任务[%s]失败: %v", task.ID.Hex(), err) // 重试逻辑:最多重试3次 update := bson.M{"$inc": bson.M{"retry_count": 1}} if task.RetryCount >= 2 { update["$set"] = bson.M{"status": "failed"} } _, _ = notificationTasksCollection.UpdateByID(ctx, task.ID, update) continue } // 标记任务为已执行 _, _ = notificationTasksCollection.UpdateByID(ctx, task.ID, bson.M{ "$set": bson.M{"status": "executed"}, }) } case <-ctx.Done(): log.Println("任务调度器已停止") return } } }
4. 手动操作的任务同步
当用户手动触发start/end路由时,需要同步更新对应任务状态,避免重复推送:
func startBooking(c *fiber.Ctx) error { bookingID, err := primitive.ObjectIDFromHex(c.Params("id")) if err != nil { return err } // 更新预订状态 _, err = bookingsCollection.UpdateByID(c.Context(), bookingID, bson.M{ "$set": bson.M{"active": true}, }) if err != nil { return err } // 标记"started"任务为已执行(避免调度器重复推送) _, _ = notificationTasksCollection.UpdateMany(c.Context(), bson.M{ "booking_id": bookingID, "type": "started", "status": "pending", }, bson.M{"$set": bson.M{"status": "executed"}}) return c.SendStatus(fiber.StatusOK) }
二、推送服务选型
1. 继续使用Expo Push API
适合React Native项目,优势是集成简单,无需维护iOS/Android推送证书,跨平台支持完善,适合快速迭代的移动端业务。
2. 备选方案
- Firebase Cloud Messaging (FCM):Google官方推送服务,支持多平台(Android/iOS/Web),免费额度充足,适合需要大规模推送的场景,生态集成性强。
- OneSignal:提供可视化推送管理界面,支持用户分段、A/B测试、推送统计等功能,适合需要精细化运营的项目,免费版足够中小规模使用。
三、关键注意事项
- 时区统一:所有时间字段必须用UTC存储和计算,避免不同时区导致推送时间错误。
- 索引优化:给
NotificationTask集合的trigger_at和status字段创建复合索引,提升查询效率。 - 任务取消:当用户取消预订时,需批量更新对应任务状态为
canceled,避免无效推送。 - 幂等性保障:每个任务执行后必须标记状态,防止服务器重启后重复推送。
内容的提问来源于stack exchange,提问作者haste
相关产品推荐
相关产品推荐

