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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 23:42:30