-
Notifications
You must be signed in to change notification settings - Fork 1
/
worker.go
85 lines (71 loc) · 2.18 KB
/
worker.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
package main
import (
"log"
"time"
"encoding/json"
"github.com/Hands-On-Restful-Web-services-with-Go/chapter9/longRunningTaskV1/models"
"github.com/streadway/amqp"
)
// Workers does the job. It holds a connection
type Workers struct {
conn *amqp.Connection
}
func (w *Workers) run() {
log.Printf("Workers are booted up and running")
channel, err := w.conn.Channel()
handleError(err, "Fetching channel failed")
defer channel.Close()
jobQueue, err := channel.QueueDeclare(
queueName, // Name of the queue
false, // Message is persisted or not
false, // Delete message when unused
false, // Exclusive
false, // No Waiting time
nil, // Extra args
)
handleError(err, "Job queue fetch failed")
messages, err := channel.Consume(
jobQueue.Name, // queue
"", // consumer
true, // auto-acknowledge
false, // exclusive
false, // no-local
false, // no-wait
nil, // args
)
go func() {
for message := range messages {
job := models.Job{}
err = json.Unmarshal(message.Body, &job)
log.Printf("Workers received a message from the queue: %s", job)
handleError(err, "Unable to load queue message")
switch job.Type {
case "A":
w.dbWork(job)
case "B":
w.callbackWork(job)
case "C":
w.emailWork(job)
}
}
}()
defer w.conn.Close()
wait := make(chan bool)
<-wait // Run long-running worker
}
func (w *Workers) dbWork(job models.Job) {
result := job.ExtraData.(map[string]interface{})
log.Printf("Worker %s: extracting data..., JOB: %s", job.Type, result)
time.Sleep(2 * time.Second)
log.Printf("Worker %s: saving data to database..., JOB: %s", job.Type, job.ID)
}
func (w *Workers) callbackWork(job models.Job) {
log.Printf("Worker %s: performing some long running process..., JOB: %s", job.Type, job.ID)
time.Sleep(10 * time.Second)
log.Printf("Worker %s: posting the data back to the given callback..., JOB: %s", job.Type, job.ID)
}
func (w *Workers) emailWork(job models.Job) {
log.Printf("Worker %s: sending the email..., JOB: %s", job.Type, job.ID)
time.Sleep(2 * time.Second)
log.Printf("Worker %s: sent the email successfully, JOB: %s", job.Type, job.ID)
}