mirror of
https://github.com/hibiken/asynq.git
synced 2025-09-19 05:17:30 +08:00
Remove per task heartbeat
This commit is contained in:
19
processor.go
19
processor.go
@@ -81,31 +81,12 @@ func (p *processor) exec() {
|
||||
task := &Task{Type: msg.Type, Payload: msg.Payload}
|
||||
p.sema <- struct{}{} // acquire token
|
||||
go func(task *Task) {
|
||||
quit := make(chan struct{}) // channel to signal heartbeat goroutine
|
||||
defer func() {
|
||||
quit <- struct{}{}
|
||||
if err := p.rdb.srem(inProgress, msg); err != nil {
|
||||
log.Printf("[SERVER ERROR] SREM failed: %v\n", err)
|
||||
}
|
||||
if err := p.rdb.clearHeartbeat(msg.ID); err != nil {
|
||||
log.Printf("[SERVER ERROR] DEL heartbeat failed: %v\n", err)
|
||||
}
|
||||
<-p.sema // release token
|
||||
}()
|
||||
// start "heartbeat" goroutine
|
||||
go func() {
|
||||
ticker := time.NewTicker(5 * time.Second)
|
||||
for {
|
||||
select {
|
||||
case <-quit:
|
||||
return
|
||||
case t := <-ticker.C:
|
||||
if err := p.rdb.heartbeat(msg.ID, t); err != nil {
|
||||
log.Printf("[ERROR] heartbeat failed for %v at %v: %v", msg.ID, t, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
err := p.handler(task) // TODO(hibiken): maybe also handle panic?
|
||||
if err != nil {
|
||||
retryTask(p.rdb, msg, err)
|
||||
|
Reference in New Issue
Block a user