-
Notifications
You must be signed in to change notification settings - Fork 1
/
sub.go
59 lines (52 loc) · 1.36 KB
/
sub.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
package main
import (
"flag"
"log"
"os"
"os/signal"
"github.com/gofrs/uuid"
stan "github.com/nats-io/stan.go"
)
func main() {
log.SetFlags(log.LstdFlags | log.Lshortfile)
var (
url = flag.String("url", stan.DefaultNatsURL, "NATS Server URLs, separated by commas")
clusterID = flag.String("cluster_id", "store3", "Cluster ID")
clientID = flag.String("client_id", "", "Client ID")
queueGroup = flag.String("queue-group", "", "Queue group ID")
)
flag.Parse()
if *clientID == "" {
*clientID = uuid.Must(uuid.NewV4()).String()
}
// Connect to NATS Streaming Server cluster
sc, err := stan.Connect(*clusterID, *clientID,
stan.NatsURL(*url),
stan.Pings(10, 5),
stan.SetConnectionLostHandler(func(_ stan.Conn, reason error) {
log.Printf("Connection lost: %v", reason)
}),
)
if err != nil {
log.Fatal(err)
}
defer sc.Close()
// Subscribe to the ECHO channel as a queue.
// Start with new messages as they come in; don't replay earlier messages.
sub, err := sc.QueueSubscribe("ECHO", *queueGroup, func(msg *stan.Msg) {
log.Printf("%10s | %s\n", msg.Subject, string(msg.Data))
}, stan.StartWithLastReceived())
if err != nil {
log.Fatal(err)
}
// Wait for Ctrl+C
doneCh := make(chan bool)
go func() {
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, os.Interrupt)
<-sigCh
sub.Unsubscribe()
doneCh <- true
}()
<-doneCh
}