Golang RabbitMQ 推送 JSON 数据并设置优先级
下面是一个简单的示例,说明如何使用 Golang 和 RabbitMQ 推送 JSON 数据,并设置优先级:
package main
import (
"encoding/json"
"fmt"
"github.com/streadway/amqp"
)
// 定义消息结构体
type Message struct {
Text string `json:"text"`
}
func main() {
// 连接到 RabbitMQ
conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
if err != nil {
panic(err)
}
defer conn.Close()
// 创建一个 channel
ch, err := conn.Channel()
if err != nil {
panic(err)
}
defer ch.Close()
// 声明一个 exchange
err = ch.ExchangeDeclare(
"my-exchange", // exchange 的名称
"direct", // exchange 的类型
true, // 是否持久化
false, // 是否自动删除
false, // 是否内部使用
false, // 是否等待确认
nil, // 其他属性
)
if err != nil {
panic(err)
}
// 声明一个 queue
q, err := ch.QueueDeclare(
"my-queue", // queue 的名称
true, // 是否持久化
false, // 是否自动删除
false, // 是否排他
false, // 是否等待确认
nil, // 其他属性
)
if err != nil {
panic(err)
}
// 将 queue 绑定到 exchange 上
err = ch.QueueBind(
q.Name, // queue 的名称
"my-routingkey", // routing key
"my-exchange", // exchange 的名称
false, // 是否等待确认
nil, // 其他属性
)
if err != nil {
panic(err)
}
// 创建一个消息
msg := Message{Text: "Hello, world!"
}
body, err := json.Marshal(msg)
if err != nil {
panic(err)
}
// 发布消息
err = ch.Publish(
"my-exchange", // exchange 的名称
"my-routingkey", // routing key
false, // 是否等待确认
false, // 是否等待确认
amqp.Publishing{ // 消息属性
ContentType: "application/json",
Body: body,
Priority: 1, // 设置优先级
},
)
if err != nil {
panic(err)
}
fmt.Println("Message sent!")
}
在上面的示例中,我们首先使用 amqp.Dial 连接到 RabbitMQ。然后,我们使用 conn.Channel 创建一个 channel。接下来,我们使用 ch.ExchangeDeclare 声明一个 exchange,使用 ch.QueueDeclare 声明一个 queue,并使用 ch.QueueBind 将 queue 绑定到 exchange 上。
我们创建一个名为 Message 的结构体,其中包含一个字符串字段 Text。然后,我们创建一个 Message 实例,并使用 json.Marshal 将其序列化为 JSON 格式的字节数组。
最后,我们使用 ch.Publish 将消息发布到 RabbitMQ。我们将 Priority 属性设置为 1,以设置消息的优先级。
原文地址: https://www.cveoy.top/t/topic/lAhl 著作权归作者所有。请勿转载和采集!