全部 API 就这六个函数
就这六个。参数开头都一样:先是流程上下文,然后一个标识这一步的 key。它们之间用普通的 if 和 for 连起来。
Task
func Task(c *Context, key string, fun TaskFun, opts ...TaskOption)执行一步。返回 error 就按退避重试,重试次数用完,整个流程失败。WithMaxRetry(n) 的 n 是重试次数,不是执行次数,所以最多执行 n+1 遍,默认 3。WithMaxRetry(0) 表示失败即失败。想主动停下而不判失败,返回 gotick.AbortError。已经失败的步骤,下一轮会中止整个流程,不会被跳过,所以坏掉的步骤不会放后面的步骤过去。
gotick.Task(ctx, "charge", func(c *gotick.TaskContext) error { return billing.Charge(c, orderId) // c is a context.Context}, gotick.WithMaxRetry(5))Memo
func Memo[T interface{}](ctx *Context, key string, build func() (T, error), opts ...TaskOption) T执行一步并记住结果。后面每一遍都直接拿回这个值,不会再调用你的构造函数。查库的结果、生成的 id、当前时刻,都靠它安全地进到流程里。流程体可以放心读这个结果,因为每一遍读到的都是同一个。
user := gotick.Memo(ctx, "user", func() (User, error) { return db.GetUser(userId)}) // safe to branch on: every run reads the same valueif user.Plan == "pro" { gotick.Task(ctx, "notify-csm", notifyCSM)}Sleep
func Sleep(c *Context, key string, duration time.Duration)等一会儿,期间不占进程。gotick 把状态记下来,把事件排到之后,这段时间你没有任何东西在跑,中途重启也不影响。醒来的精度取决于 DelayedTaskCheckInterval,默认 500 毫秒。
gotick.Sleep(ctx, "wait-payment", 30*time.Minute)WaitForSignal
func WaitForSignal[T any](c *Context, key string, opts ...SignalOption) (T, bool)等外面的事情发生。返回的 bool 告诉你信号来了没有,false 就是超时赢了。不加 WithSignalTimeout 就是无限等,这种情况不会安排任何唤醒事件,只有 Cancel 能把它捞出来。第一个信号说了算,后来的一律拒掉,所以每一遍重放都读到同一个值。比等待点更早到的信号会先存着,不会丢。
paid, ok := gotick.WaitForSignal[Payment](ctx, "paid", gotick.WithSignalTimeout(30*time.Minute)) if !ok { gotick.Task(ctx, "close-order", closeOrder) return}gotick.Task(ctx, "ship", func(c *gotick.TaskContext) error { return shipping.Send(c, paid.TradeNo)})Array
func Array[T interface{}](ctx *Context, key string, build func(ctx *TaskContext) ([]T, error), opts ...TaskOption) []ArrayWrap[T]列表版的 Memo。后面有步骤要遍历这个结果时就该用它:列表存下来了,每一遍循环的长度和元素都一样。Value() 取出一个元素,Key(prefix) 给这个元素单独一个步骤 key。
files := gotick.Array(ctx, "list", func(c *gotick.TaskContext) ([]string, error) { return storage.List(c, dir)}) for _, f := range files { name := f.Value() gotick.Task(ctx, f.Key("convert"), func(c *gotick.TaskContext) error { return convert(c, name) })}Async + Wait
func Async[T interface{}](ctx *Context, key string, f func(ctx *TaskContext) (T, error), opts ...TaskOption) *FutureT[T]func AsyncArray[T, A interface{}](ctx *Context, key string, arr []ArrayWrap[A], f func(ctx *TaskContext, a A, index int) (T, error), opts ...TaskOption) []Futurefunc Wait(ctx *Context, parallel int, fs ...Future)声明几个可以同时跑的步骤,再用 Wait 限并发地等它们。AsyncArray 是「按 Array 展开成并行任务」的简写。Wait 返回之后用 Value() 读结果。WithMaxRetry 对它们同样有效,只是不传时的默认值是 5 次重试,而 Task 是 3 次。
这里有个别扭的地方:AsyncArray 返回 []Future,而 Future 上没有 Value(),取结果得先类型断言回具体的 future 类型。单独用 Async 返回的是 *FutureT[T],不用断言。
fs := gotick.AsyncArray(ctx, "download", files, func(c *gotick.TaskContext, url string, i int) (string, error) { return download(c, url) }) gotick.Wait(ctx, 4, fs...) // at most 4 at a time gotick.Task(ctx, "save", func(c *gotick.TaskContext) error { for _, f := range fs { path := f.(*gotick.FutureT[string]).Value() _ = path } return nil})Sequence
func Sequence(ctx *Context, key string, maxLen int) SequenceWrap一个可恢复的计数器,用在长度不是列表的循环上。Next() 往前走一步并存下位置,TaskKey(prefix) 给这一轮一个独立的步骤 key。maxLen 传负数就一直循环,由你自己 break。
seq := gotick.Sequence(ctx, "pages", 100) for seq.Next() { page := seq.Current gotick.Task(ctx, seq.TaskKey("fetch"), func(c *gotick.TaskContext) error { return fetchPage(c, page) })}流程结束时
Flow 上挂着三个回调,它们本身也是普通步骤。OnSuccess 在函数正常返回之后跑一次。OnFail 在流程对某一步彻底放弃时跑。OnError 每次单步失败都跑,包括之后还会重试的那些。所以报警用 OnError,补偿用 OnFail。
tick.Flow("order/close", closeOrderFlow). OnSuccess(func(ctx *gotick.Context) error { return metrics.Inc("order.closed") }). OnError(func(ctx *gotick.Context, ts gotick.TaskStatus) error { return alert.Warn(ctx.CallId, ts.Errs) // fires on every attempt }). OnFail(func(ctx *gotick.Context, ts gotick.TaskStatus) error { return compensate(ctx.CallId) // fires once, after giving up })