当前位置: 首页
编程语言
Linux下Golang实现消息队列的完整指南

Linux下Golang实现消息队列的完整指南

热心网友 时间:2026-07-20
转载

在Linux下用Golang实现消息队列,常用RabbitMQ、Kafka、Redis。RabbitMQ支持多协议、生态成熟,可通过包管理器安装,借助streadway amqp库实现生产者和消费者;Kafka高吞吐,Redis轻量,各具优势。

在Linux环境下用Golang实现消息队列,其实有不少选择。最常用的三种方案——RabbitMQ、Kafka、Redis——各有各的适用场景,下面一个个拆开来看。

Linux中Golang如何实现消息队列

1. 使用RabbitMQ

RabbitMQ是消息队列领域的“老牌选手”,支持多种协议,生态成熟,部署起来也相对简单。

安装RabbitMQ

在Linux上安装,直接走包管理器就行:

sudo apt-get update
sudo apt-get install rabbitmq-server

使用Golang客户端库

Go这边最常用的库是 streadway/amqp,安装命令:

go get github.com/streadway/amqp

示例代码

先看一个生产者,负责往队列里发消息:

package main

import (
    "fmt"
    "log"
    "github.com/streadway/amqp"
)

func failOnError(err error, msg string) {
    if err != nil {
        log.Fatalf("%s: %s", msg, err)
    }
}

func main() {
    conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
    failOnError(err, "Failed to connect to RabbitMQ")
    defer conn.Close()

    ch, err := conn.Channel()
    failOnError(err, "Failed to open a channel")
    defer ch.Close()

    q, err := ch.QueueDeclare(
        "hello", // name
        true,    // durable
        false,   // delete when unused
        false,   // exclusive
        false,   // no-wait
        nil,     // arguments
    )
    failOnError(err, "Failed to declare a queue")

    body := "Hello World!"
    err = ch.Publish(
        "",     // exchange
        q.Name, // routing key
        false,  // mandatory
        false,  // immediate
        amqp.Publishing{
            ContentType: "text/plain",
            Body:        []byte(body),
        })
    failOnError(err, "Failed to publish a message")
    fmt.Println(" [x] Sent %s", body)
}

消费者这边,逻辑也很清晰——监听队列,收到消息就处理:

package main

import (
    "fmt"
    "log"
    "github.com/streadway/amqp"
)

func failOnError(err error, msg string) {
    if err != nil {
        log.Fatalf("%s: %s", msg, err)
    }
}

func main() {
    conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
    failOnError(err, "Failed to connect to RabbitMQ")
    defer conn.Close()

    ch, err := conn.Channel()
    failOnError(err, "Failed to open a channel")
    defer ch.Close()

    q, err := ch.QueueDeclare(
        "hello", // name
        true,    // durable
        false,   // delete when unused
        false,   // exclusive
        false,   // no-wait
        nil,     // arguments
    )
    failOnError(err, "Failed to declare a queue")

    msgs, err := ch.Consume(
        q.Name, // queue
        "",     // consumer
        true,   // auto-ack
        false,  // exclusive
        false,  // no-local
        false,  // no-wait
        nil,    // args
    )
    failOnError(err, "Failed to register a consumer")

    forever := make(chan bool)

    go func() {
        for d := range msgs {
            fmt.Printf("Received a message: %s\n", d.Body)
        }
    }()

    fmt.Println(" [*] Waiting for messages. To exit press CTRL+C")
    <-forever
}

2. 使用Kafka

如果你需要处理海量数据流,或者对高吞吐量有刚需,Kafka是更合适的选择。它本质上是一个分布式流处理平台,但做消息队列也完全够用。

安装Kafka

下载解压,然后启动Zookeeper和Kafka服务器:

wget https://downloads.apache.org/kafka/2.8.0/kafka_2.13-2.8.0.tgz
tar -xzf kafka_2.13-2.8.0.tgz
cd kafka_2.13-2.8.0
bin/zookeeper-server-start.sh config/zookeeper.properties &
bin/kafka-server-start.sh config/server.properties &

使用Golang客户端库

推荐用 confluent-kafka-go,安装:

go get github.com/confluentinc/confluent-kafka-go/kafka

示例代码

生产者示例:

package main

import (
    "fmt"
    "github.com/confluentinc/confluent-kafka-go/kafka"
)

func main() {
    p, err := kafka.NewProducer(&kafka.ConfigMap{
        "bootstrap.servers": "localhost:9092"})
    if err != nil {
        panic(err)
    }
    defer p.Close()

    go func() {
        for e := range p.Events() {
            switch ev := e.(type) {
            case *kafka.Message:
                if ev.TopicPartition.Error != nil {
                    fmt.Printf("Delivery failed: %v\n", ev.TopicPartition.Error)
                } else {
                    fmt.Printf("Delivered message to %v\n", ev.TopicPartition)
                }
            }
        }
    }()

    topic := "test-topic"
    p.Produce(&kafka.Message{
        TopicPartition: kafka.TopicPartition{
            Topic: &topic, Partition: kafka.PartitionAny},
        Value: []byte("Hello Kafka"),
    }, nil)

    p.Flush(15 * 1000)
}

消费者示例:

package main

import (
    "fmt"
    "github.com/confluentinc/confluent-kafka-go/kafka"
)

func main() {
    c, err := kafka.NewConsumer(&kafka.ConfigMap{
        "bootstrap.servers": "localhost:9092",
        "group.id":          "test-consumer-group",
        "auto.offset.reset": "earliest",
    })
    if err != nil {
        panic(err)
    }
    defer c.Close()

    c.SubscribeTopics([]string{"test-topic"}, nil)

    for {
        msg, err := c.ReadMessage(-1)
        if err == nil {
            fmt.Printf("Received message: %s\n", string(msg.Value))
        } else {
            fmt.Printf("Consumer error: %v\n", err)
        }
    }
}

3. 使用Redis

Redis虽然以缓存闻名,但它的发布/订阅功能实现一个轻量级消息队列绰绰有余,特别适合对实时性要求不高、系统复杂度想控制得低一些的场景。

安装Redis

sudo apt-get update
sudo apt-get install redis-server

使用Golang客户端库

go-redis 库,安装:

go get github.com/go-redis/redis/v8

示例代码

生产者:

package main

import (
    "context"
    "fmt"
    "github.com/go-redis/redis/v8"
)

var ctx = context.Background()

func main() {
    rdb := redis.NewClient(&redis.Options{
        Addr:     "localhost:6379",
        Password: "", // no password set
        DB:       0,  // use default DB
    })

    err := rdb.Publish(ctx, "channel", "Hello Redis").Err()
    if err != nil {
        panic(err)
    }
    fmt.Println("Message published")
}

消费者:

package main

import (
    "context"
    "fmt"
    "github.com/go-redis/redis/v8"
)

var ctx = context.Background()

func main() {
    rdb := redis.NewClient(&redis.Options{
        Addr:     "localhost:6379",
        Password: "",
        DB:       0,
    })

    pubsub := rdb.Subscribe(ctx, "channel")
    defer pubsub.Close()

    ch := pubsub.Channel()
    for msg := range ch {
        fmt.Printf("Received message: %s\n", msg.Payload)
    }
}

这三种方案各有侧重:RabbitMQ胜在协议丰富、功能完善;Kafka强在吞吐量和持久化;Redis胜在简单、部署成本低。实际项目里怎么选,取决于你的业务场景——是追求可靠性、吞吐量,还是图个轻量快速。上面这些代码直接跑起来就能用,剩下的就看你的具体需求了。

来源:https://www.yisu.com/ask/70065612.html

游乐网为非赢利性网站,所展示的游戏/软件/文章内容均来自于互联网或第三方用户上传分享,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系youleyoucom@outlook.com。

同类文章
更多
CentOS系统PHP超时问题详细原因分析与全面解决方法

CentOS系统PHP超时问题详细原因分析与全面解决方法

CentOS系统PHP超时问题源于脚本执行超限,可从配置参数、代码性能、服务器环境三方面解决:修改php ini或 htaccess全局设置,脚本内动态调整超时,同步优化Web服务器超时参数,并排查慢查询、外部依赖等性能瓶颈。

时间:2026-07-20 21:33
深入分析Java日志对CentOS性能的影响

深入分析Java日志对CentOS性能的影响

Java日志对CentOS性能的影响主要体现在I O、CPU和内存占用。频繁写入磁盘增加I O负载,异步日志或缓冲区可缓解;日志级别过低导致CPU占用升高,需合理设置;缓冲区过大会增加内存压力,可调整大小或使用内存映射文件。日志存储需轮转压缩归档,分析处理宜采用自动化工具。

时间:2026-07-20 21:33
CentOS Cobbler集成其他工具实践指南

CentOS Cobbler集成其他工具实践指南

Cobbler利用Kickstart脚本集成Puppet、Ansible等配置管理工具,实现系统安装后自动化部署;与OpenStack协同管理虚拟机镜像及创建流程;配合GlusterFS自动挂载存储卷;并原生支持DNS、DHCP网络自动配置,全面提升部署效率。

时间:2026-07-20 21:33
CentOS上C++性能优化选项配置完整指南

CentOS上C++性能优化选项配置完整指南

在CentOS上配置C++性能优化,需安装gcc等编译工具,使用-O2、-march=native、-mtune=native及-flto编译选项,借助gprof、perf、valgrind定位性能热点,同时可通过环境变量或CMake统一管理优化设置,推荐结合CPU架构特性并利用CMake条件编译,从而显著提高程序执行速度与效率,实现最佳性能。

时间:2026-07-20 21:33
Crontab时间格式错误如何修复

Crontab时间格式错误如何修复

Crontab时间格式由5个字段组成,分别代表分钟、小时、日期、月份和星期。常见错误包括多余字符、字段值超出范围、缺少字段、星期值错误以及字段间缺少空格或制表符。对照这些典型问题逐项排查即可解决。

时间:2026-07-20 06:48
热门专题
更多
刀塔传奇破解版无限钻石下载大全 刀塔传奇破解版无限钻石下载大全
洛克王国正式正版手游下载安装大全 洛克王国正式正版手游下载安装大全
思美人手游下载专区 思美人手游下载专区
好玩的阿拉德之怒游戏下载合集 好玩的阿拉德之怒游戏下载合集
不思议迷宫手游下载合集 不思议迷宫手游下载合集
百宝袋汉化组游戏最新合集 百宝袋汉化组游戏最新合集
jsk游戏合集30款游戏大全 jsk游戏合集30款游戏大全
宾果消消消原版下载大全 宾果消消消原版下载大全
  • 热门数据榜