首页 >后端开发 >Golang >用 Go 构建 Kafka 生产者和消费者

用 Go 构建 Kafka 生产者和消费者

Barbara Streisand
Barbara Streisand原创
2025-01-03 19:48:43827浏览

Building a Kafka Producer and Consumer in Go

Apache Kafka 是一个强大的分布式流平台,用于构建实时数据管道和流应用程序。在这篇博文中,我们将逐步使用 Golang 设置 Kafka 生产者和消费者。

先决条件

在我们开始之前,请确保您的计算机上安装了以下软件:

  • Go(1.16 或更高)

  • Docker(用于在本地运行 Kafka)

  • 卡夫卡

使用 Docker 设置 Kafka

为了快速设置 Kafka,我们将使用 Docker。在项目目录中创建 docker-compose.yml 文件:

yamlCopy codeversion: '3.7'
services:
  zookeeper:
    image: wurstmeister/zookeeper:3.4.6
    ports:
      - "2181:2181"
  kafka:
    image: wurstmeister/kafka:2.13-2.7.0
    ports:
      - "9092:9092"
    environment:
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
    depends_on:
      - zookeeper

运行以下命令启动 Kafka 和 Zookeeper:

docker-compose up -d

在 Go 中创建 Kafka 生产者

首先,初始化一个新的 Go 模块:

go mod init kafka-example

安装 kafka-go 库:

go get github.com/segmentio/kafka-go

现在,创建一个文件 Producer.go 并添加以下代码:

package main

import (
    "context"
    "fmt"
    "github.com/segmentio/kafka-go"
    "log"
    "time"
)

func main() {
    writer := kafka.Writer{
        Addr:     kafka.TCP("localhost:9092"),
        Topic:    "example-topic",
        Balancer: &kafka.LeastBytes{},
    }

    defer writer.Close()

    for i := 0; i < 10; i++ {
        msg := kafka.Message{
            Key:   []byte(fmt.Sprintf("Key-%d", i)),
            Value: []byte(fmt.Sprintf("Hello Kafka %d", i)),
        }

        err := writer.WriteMessages(context.Background(), msg)
        if err != nil {
            log.Fatal("could not write message " + err.Error())
        }

        time.Sleep(1 * time.Second)
        fmt.Printf("Produced message: %s\n", msg.Value)
    }
}

此代码设置一个 Kafka 生产者,向 example-topic 主题发送 10 条消息。

运行生产者:

go run producer.go

您应该看到指示消息已生成的输出。

在 Go 中创建 Kafka 消费者

创建文件consumer.go并添加以下代码:

package main

import (
    "context"
    "fmt"
    "github.com/segmentio/kafka-go"
    "log"
)

func main() {
    reader := kafka.NewReader(kafka.ReaderConfig{
        Brokers: []string{"localhost:9092"},
        Topic:   "example-topic",
        GroupID: "example-group",
    })

    defer reader.Close()

    for {
        msg, err := reader.ReadMessage(context.Background())
        if err != nil {
            log.Fatal("could not read message " + err.Error())
        }
        fmt.Printf("Consumed message: %s\n", msg.Value)
    }
}

该消费者从 example-topic 主题读取消息并将其打印到控制台。

运行消费者:

go run consumer.go

您应该看到指示消息已被消耗的输出。

结论

在这篇博文中,我们演示了如何使用 Golang 设置 Kafka 生产者和消费者。这个简单的示例展示了生成和消费消息的基础知识,但 Kafka 的功能远远不止于此。借助 Kafka,您可以构建强大的、可扩展的实时数据处理系统。

随意探索更高级的功能,例如消息分区、基于密钥的消息分发以及与其他系统的集成。快乐编码!


就是这样!这篇博文简要介绍了如何将 Kafka 与 Go 结合使用,非常适合想要开始实时数据处理的开发人员。

以上是用 Go 构建 Kafka 生产者和消费者的详细内容。更多信息请关注PHP中文网其他相关文章!

声明:
本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系admin@php.cn