美文网首页
基于 RabbitMQ 实现数据异步入库

基于 RabbitMQ 实现数据异步入库

作者: WuCh1k1n | 来源:发表于2020-01-26 12:14 被阅读0次

    RabbitMQ 简介

    AMQP,即Advanced Message Queuing Protocol,高级消息队列协议,是应用层协议的一个开放标准,为面向消息的中间件设计。消息中间件主要用于组件之间的解耦,消息的发送者无需知道消息使用者的存在,反之亦然。
    AMQP的主要特征是面向消息、队列、路由(包括点对点和发布/订阅)、可靠性、安全。
    RabbitMQ是一个开源的AMQP实现,服务器端用Erlang语言编写,支持多种客户端,如:Python、Ruby、.NET、Java、JMS、C、PHP、ActionScript、XMPP、STOMP等,支持AJAX。用于在分布式系统中存储转发消息,在易用性、扩展性、高可用性等方面表现不俗。


    消息队列异步入库

    那么,基于 RabbitMQ(消息队列)实现数据异步入库有什么好处呢?

    1. 通过异步处理提高系统性能(减少响应时间,降低数据库压力)
    2. 降低系统耦合性
      消息队列异步入库.png
      在不使用消息队列服务器的时候,应用会直接将用户请求导致的数据变更写入数据库,在高并发的情况下数据库的写入压力会剧增,系统整体响应速度明显变慢。
      但是在使用消息队列之后,用户请求导致的数据变更发送给消息队列之后立即返回,在高并发的情况下也不会增加数据库的写入压力,系统整体响应速度无明显变化。再由消息队列的消费者进程从消息队列中获取数据,异步写入数据库。由于消息队列服务器处理速度快于数据库(消息队列也比数据库有更好的伸缩性),因此响应速度得到大幅改善。
      因此,消息队列具备良好的削峰功能,将短时间高并发产的事务数据存储在消息队列之中,再异步地平缓写入数据库之中,减轻数据库的并发写入压力,提高系统整体的访问性能。消息队列的应用常见于各种稀缺资源抢购系统。

    消息队列结构体

    RabbitMQ客户端模块

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

    消息队列生产端

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

    消息队列消费端

    // simple模式的消费端
    func (r *RabbitMQ) ConsumeSimple(orderService services.IOrderService, productService services.IProductService) {
        // 1.申请队列,如果队列不存在会自动创建,存在则跳过创建
        q, err := r.channel.QueueDeclare(
            r.QueueName,
            // 是否持久化
            false,
            // 是否自动删除
            false,
            // 是否具有排他性
            false,
            // 是否阻塞处理
            false,
            // 额外的属性
            nil,
        )
        if err != nil {
            fmt.Println(err)
        }
    
        // 消费者流控
        r.channel.Qos(
            1, // 当前消费者一次能接受的最大消息数量
            0, // 服务器传递的最大容量(以八位字节为单位)
            false, // 如果设置为true对channel可用
        )
    
        // 接收消息
        msgs, err := r.channel.Consume(
            q.Name, // queue
            "", // consumer,用来区分多个消费者
            false, // auto-ack,是否自动应答
            false, // exclusive,是否独有
            false, // no-local,设置为true表示不能将同一个Connection中生产者发送的消息传递给这个Connection中的消费者
            false, // no-wait,队列是否阻塞
            nil,   // args
        )
        if err != nil {
            fmt.Println(err)
        }
    
        forever := make(chan bool)
        //启用协程处理消息
        go func() {
            for d := range msgs {
                // 消息逻辑处理,可以自行设计逻辑
                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 CTRL+C")
        <-forever
    }
    

    参考:
    RabbitMQ基础概念详细介绍
    消息队列

    相关文章

      网友评论

          本文标题:基于 RabbitMQ 实现数据异步入库

          本文链接:https://www.haomeiwen.com/subject/izkrthtx.html