-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathrabbitmq.go
More file actions
120 lines (104 loc) · 3.17 KB
/
Copy pathrabbitmq.go
File metadata and controls
120 lines (104 loc) · 3.17 KB
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
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
package event
import (
"fmt"
"strings"
"github.com/assembla/cony"
"github.com/streadway/amqp"
)
const (
// TODO: delayQueue 与 deadQueue 都需要与 hutch, hutch-schedule 命名同规则, 这样才可以整合使用
deadQueue = "dead_queue"
delayQueue = "delay_queue"
exchangeName = "hutch"
//# fixed delay levels
//# seconds(4): 5s, 10s, 20s, 30s
//# minutes(14): 1m, 2m, 3m, 4m, 5m, 6m, 7m, 8m, 9m, 10m, 20m, 30m, 40m, 50m
//# hours(3): 1h, 2h, 3h
)
var (
// 拥有的 delayQueues
delayQueues = []string{
"5s", "10s", "20s", "30s",
"60s", "120s", "180s", "240s", "300s", "360s", "420s", "480s", "540s", "600s", "1200s", "1800s", "2400s", "3000s",
"3600s", "7200s", "10800s",
}
// 30 * 24 * 3600 * 1000
oneMonth = int64(2592000000)
)
// 构建 buildin 重试/retry 的队列
func buildinQueueName(names ...string) string {
return strings.Join(names, "_")
}
// 构建内置的 retry 与 dead queue
func builtinQueue(exc cony.Exchange, delayExc cony.Exchange, cli *cony.Client) {
var declears []cony.Declaration
// delay queue_<5s>: 不需要 consumer, 由 rabbitmq 的 ddl 自行处理
// - 最长超过 30 天重新投递
// - 重新投递到默认的 exchange
qus, binds := builtinDelayQueues(exc, delayExc)
// dead queue: 不需要 consumer, 由 rabbitmq 自行过期处理
// - 消息超过 30 天放弃
// - 超过 10w 条消息放弃
deadQue := buildConyQueue(
buildinQueueName(exc.Name, deadQueue),
amqp.Table{"x-message-ttl": oneMonth, "x-max-length": int64(100000)},
)
deadBnd := cony.Binding{Queue: deadQue, Exchange: exc, Key: deadQueue + ".#"}
qus = append(qus, deadQue)
binds = append(binds, deadBnd)
for _, q := range qus {
declears = append(declears, cony.DeclareQueue(q))
}
for _, b := range binds {
declears = append(declears, cony.DeclareBinding(b))
}
cli.Declare(declears)
}
func builtinDelayQueues(exc cony.Exchange, delayExc cony.Exchange) ([]*cony.Queue, []cony.Binding) {
var ques []*cony.Queue
var binds []cony.Binding
for _, q := range delayQueues {
retryQue := buildConyQueue(
// <hutch>_delay_queue_<5s>
buildinQueueName(exc.Name, delayQueue, q),
amqp.Table{"x-message-ttl": oneMonth, "x-dead-letter-exchange": exc.Name},
)
ques = append(ques, retryQue)
binds = append(binds, cony.Binding{
Queue: retryQue,
Exchange: delayExc,
Key: fmt.Sprintf("%s.schedule.%s", exc.Name, q), // <hutch>.schedule.<5s>
})
}
return ques, binds
}
// 构建一个默认的持久化的 queue
func buildConyQueue(name string, table amqp.Table) *cony.Queue {
return &cony.Queue{
Name: name,
AutoDelete: false,
Durable: true,
Args: table,
}
}
// buildDefaultOpt 构建默认的 cony.Client 配置
func buildDefaultOpt() []cony.ClientOpt {
return []cony.ClientOpt{
cony.URL(""),
cony.Backoff(cony.DefaultBackoff),
}
}
func buildScheduleExchange(name string) cony.Exchange {
vname := name
if len(name) == 0 {
vname = exchangeName
}
return buildTopicExchange(fmt.Sprintf("%s.schedule", vname))
}
func buildTopicExchange(name string) cony.Exchange {
vname := name
if len(name) == 0 {
vname = exchangeName
}
return cony.Exchange{Name: vname, AutoDelete: false, Durable: true, Kind: "topic"}
}