GO语言实现RabbitMQ的几种工作模式

我们利用Go语言来实现RabbitMQ的简单模式,其他的工作模式是根据简单模式进行构建的。

RabbitMQ的实例代码如下所示:

const MQURL = "amqp://用户名:密码@host:端口/vhost"

type RabbitMQ struct {
    conn *amqp.Connection
    channel *amqp.Channel
    //队列名称
    QueueName string
    //交换机
    Exchange string
    //key
    Key string
    //连接信息
     MqUrl string
    sync.Mutex
}

func NewRabbitMQ(queueName string, exchange string, key string) *RabbitMQ{
    rabbitMQ := &RabbitMQ{
        QueueName: queueName,
        Exchange:  exchange,
        Key:       key,
        MqUrl:MQURL,
    }
    var err error
    rabbitMQ.conn,err = amqp.Dial(rabbitMQ.MqUrl)
    if err != nil {
        rabbitMQ.failOnErr(err,"创建连接错误")
    }
    rabbitMQ.channel,err = rabbitMQ.conn.Channel()
    if err != nil {
        rabbitMQ.failOnErr(err,"获取Channel失败")
    }
    return rabbitMQ
}

//断开channel conn
func (r *RabbitMQ) Destroy(){
    r.channel.Close()
    r.conn.Close()
}

//错误处理函数
func (r *RabbitMQ) failOnErr(err error,msg string) {
    if err != nil {
        log.Fatal("%s:%s",msg,err)
        panic(fmt.Sprintf("%s:%s",msg,err))
    }
}

简单模式

我们来实现以下生产端的代码:

//实现simple模式
func NewRabbitMQSimple(queueName string) *RabbitMQ {
    return NewRabbitMQ(queueName,"","")
}

//实现简单模式下生产者代码
func (r *RabbitMQ) PublishSimple(message string) error{
    r.Lock()
    defer r.Unlock()
    //申请队列,如果不存在会自动创建,存在跳过创建,保证队列存在,消息能发送到队列中
    _,err := r.channel.QueueDeclare(
            r.QueueName,
            //控制消息是否持久化,true开启
            false,
            //是否为自动删除
            false,
            //是否具有排他性
            false,
            //是否阻塞
            false,
            //额外属性
            nil,
        )
    if err != nil{
        return err
    }
    //发送消息到队列中
    r.channel.Publish(
            r.Exchange,
            r.QueueName,
            //如果为true,根据exchange类型和routekey类型,如果无法找到符合条件的队列,name会把发送的信息返回给发送者
            false,
            //如果为true,当exchange发送到消息队列后发现队列上没有绑定的消费者,则会将消息返还给发送者
            false,
                        //发送信息
            amqp.Publishing{
                ContentType:     "text/plain",
                Body:            []byte(message),
            },
        )
    return nil
}

对于高并发来说,防止资源产生竞争关系,我在生成端添加了互斥锁,要保证信息已经成功的被添加到队列中;队列创建好,我们就把产生的消息发送出去。相关参数介绍我已经在代码里做了备注。

接下来我们来实现,消费端的代码:

//实现简单模式下消费者代码
func (r *RabbitMQ) ConsumeSimple(orderService service.IOrderService,
    productService service.ProductService){
    //申请队列,如果不存在会自动创建,存在跳过创建,保证队列存在,消息能发送到队列中
    q,err := r.channel.QueueDeclare(
        r.QueueName,
        //控制消息是否持久化,true开启
        false,
        //是否为自动删除
        false,
        //是否具有排他性
        false,
        //是否阻塞
        false,
        //额外属性
        nil,
    )
    if err != nil{
        fmt.Println(err)
    }

    //消费者流控,防止暴库
    r.channel.Qos(
            //每次只接受一个消息进行消费
            1,
            //服务器传递的最大容量(以8位字节为单位)
            0,
            //true对全局可用,false只对当前channel可用
            false,
        )

    //接收消息
    messages,err := r.channel.Consume(
            q.Name,
            //用来区分多个消费者
            "",
            //是否自动应答
            false,
            //是否具有排他性
            false,
            //如果设置为true,表示不能将同一个connection中发送的消息
            //传递给同一个connection的消费者
            false,
            //是否为阻塞
            false,
            nil,
        )
    if err != nil {
        fmt.Println(err)
    }
    forever := make(chan bool)
    //启用协程处理消息
    go func() {
        for d := range messages{
            //实现我们的处理逻辑函数
            log.Printf("Received a message : %s",d.Body)

            message := &datamodels.Message{}
            err :=json.Unmarshal([]byte(d.Body),message)
            if err !=nil {
                fmt.Println(err)
            }
            //插入订单
            _,err = orderService.InsertOrderByMessage(message)
            if err != nil {
                fmt.Println(err)
            }
            //扣除数量
            err = productService.SubNumberOne(message.ProductID)
            if err != nil{
                fmt.Println(err)
            }
            //如果为true表示确认所有未确认的消息,false为当前消息
            d.Ack(false)
        }
    }()

    log.Printf("[*] Waiting for messages,To exit press CTRAL+C")
    <-forever
}

消费端同样首先判断队列是否生成成功,接下来要注意的地方就是,RabbitMQ的限流,等消费者消费完一个信息,才去接收下一个信息,防止消费端在处理类似于Mysql数据库的时候发生暴库,具体实现如下:

//消费者流控,防止暴库
    r.channel.Qos(
            //每次只接受一个消息进行消费
            1,
            //服务器传递的最大容量(以8位字节为单位)
            0,
            //true对全局可用,false只对当前channel可用
            false,
        )

接下来就开始接收和消费信息,我们通过Go语言的协程来处理接收到的信息,代码如下:

forever := make(chan bool)
    //启用协程处理消息
    go func() {
        for d := range messages{
            //实现我们的处理逻辑函数
            log.Printf("Received a message : %s",d.Body)

            message := &datamodels.Message{}
            err :=json.Unmarshal([]byte(d.Body),message)
            if err !=nil {
                fmt.Println(err)
            }
            //插入订单
            _,err = orderService.InsertOrderByMessage(message)
            if err != nil {
                fmt.Println(err)
            }
            //扣除数量
            err = productService.SubNumberOne(message.ProductID)
            if err != nil{
                fmt.Println(err)
            }
            //如果为true表示确认所有未确认的消息,false为当前消息
            d.Ack(false)
        }
    }()

    log.Printf("[*] Waiting for messages,To exit press CTRAL+C")
    <-forever

注意我们将消费端自动应答模式关闭,是为了确保消费者业务逻辑处理成功再接收下一条,这里最关键的是要手动应答d.Ack(false)

工作模式

工作模式对于简单模式的代码来说是一样的,只是消费端的个数增加了,这种情况适用于生产者生成消息过于快,消费端来不及消费,才会采用多个消费端进行消费。

订阅模式

实例订阅模式的代码如下:

//实现订阅模式
func NewRabbitMQSimple(exchangeName string) *RabbitMQ {
    return NewRabbitMQ("",exchangeName ,"")
}

订阅模式下的生产端的代码如下(注意:这块和简单模式的区别在于,要生成交换机,而不是队列):

//实现订阅模式下生产者代码
func (r *RabbitMQ) PublishPub(message string) error{
    r.Lock()
    defer r.Unlock()
    err := r.channel.ExchangeDeclare(
            r.Exchange,
            "fanout", //设置交换机的类型,在订阅模式下交换机的类型为广播类型
            true,
            false,
            false,//true表示这个参数exchange不可以被client用来推送信息,仅用来exchange和exchange之间的绑定
            false,
            nil,
        )
    if err != nil {
        r.failOnErr(err,"创建交换机异常")
    }
    //发送消息到队列中
    r.channel.Publish(
        r.Exchange,
        "",
        //如果为true,根据exchange类型和routekey类型,如果无法找到符合条件的队列,name会把发送的信息返回给发送者
        false,
        //如果为true,当exchange发送到消息队列后发现队列上没有绑定的消费者,则会将消息返还给发送者
        false,
        amqp.Publishing{
            ContentType:     "text/plain",
            Body:            []byte(message),
        },
    )
    return nil
}

订阅模式的消费端代码:

//订阅模式下的消费端代码
func (r *RabbitMQ) ConsumePub(){
    err := r.channel.ExchangeDeclare(
        r.Exchange,
        "fanout", //设置交换机的类型,在订阅模式下交换机的类型为广播类型
        true,
        false,
        false,//true表示这个参数exchange不可以被client用来推送信息,仅用来exchange和exchange之间的绑定
        false,
        nil,
    )
    if err != nil {
        r.failOnErr(err,"创建交换机异常")
    }
    //申请队列,如果不存在会自动创建,存在跳过创建,保证队列存在,消息能发送到队列中
    q,err := r.channel.QueueDeclare(
        "", //随机生成队列名称
        //控制消息是否持久化,true开启
        false,
        //是否为自动删除
        false,
        //是否具有排他性
        true,
        //是否阻塞
        false,
        //额外属性
        nil,
    )
    if err != nil{
        r.failOnErr(err,"生产队列异常")
    }
    //绑定队列到Exchange中
    err = r.channel.QueueBind(
            q.Name,
            "",//在订阅模式下这个参数要为空,
            r.Exchange,
            false,
            nil,
        )
    //接收消息
    messages,err := r.channel.Consume(
        q.Name,
        //用来区分多个消费者
        "",
        //是否自动应答
        true,
        //是否具有排他性
        false,
        //如果设置为true,表示不能将同一个connection中发送的消息
        //传递给同一个connection的消费者
        false,
        //是否为阻塞
        false,
        nil,
    )
    if err != nil {
        fmt.Println(err)
    }
    forever := make(chan bool)
    //启用协程处理消息
    go func() {
        for d := range messages{
            //实现我们的处理逻辑函数
            log.Printf("Received a message : %s",d.Body)
        }
    }()

    log.Printf("[*] Waiting for messages,To exit press CTRAL+C")
    <-forever
}

路由模式

路由模式的实例代码如下:

//实现simple模式
func NewRabbitMQSimple(exchangeName string, routingKey string) *RabbitMQ {
    return NewRabbitMQ("",exchangeName ,routingKey )
}

注意路由模式和订阅模式的区别在于,将广播模式改为direct模式

err := r.channel.ExchangeDeclare(
            r.Exchange,
            "direct", 
            true,
            false,
            false,//true表示这个参数exchange不可以被client用来推送信息,仅用来exchange和exchange之间的绑定
            false,
            nil,
        )

在生产消息的时候区别在于:

r.channel.Publish(
        r.Exchange,
        r.Key,
        //如果为true,根据exchange类型和routekey类型,如果无法找到符合条件的队列,name会把发送的信息返回给发送者
        false,
        //如果为true,当exchange发送到消息队列后发现队列上没有绑定的消费者,则会将消息返还给发送者
        false,
        amqp.Publishing{
            ContentType:     "text/plain",
            Body:            []byte(message),
        },
    )

消费端代码实现:

在声明交换机的时候,将广播模式改为direct模式。绑定队列的时候,要将key参数进行添加,代码如下所示:

err = r.channel.QueueBind(
            q.Name,
            r.Key,
            r.Exchange,
            false,
            nil,
        )

话题模式

和路由模式的区别在于,将交换机的模式改变为topic模式,其他的不变。

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 216,651评论 6 501
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 92,468评论 3 392
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 162,931评论 0 353
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,218评论 1 292
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,234评论 6 388
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,198评论 1 299
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,084评论 3 418
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,926评论 0 274
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,341评论 1 311
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,563评论 2 333
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,731评论 1 348
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,430评论 5 343
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 41,036评论 3 326
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,676评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,829评论 1 269
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,743评论 2 368
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,629评论 2 354

推荐阅读更多精彩内容