使用golang连接kafka

人生之路坎坎坷坷,跌跌撞撞在所难免。但是,不论跌了多少次,你都必须坚强勇敢地站起来。任何时候,无论你面临着生命的何等困惑抑或经受着多少挫折,无论道路多艰难,希望变得如何渺茫,请你不要绝望,再试一次,坚持到底,成功终将属于勇不言败的你。

导读:本篇文章讲解 使用golang连接kafka,希望对大家有帮助,欢迎收藏,转发!站点地址:www.bmabk.com,来源:原文

1 下载,配置,启动 kafka

下载链接

配置修改

在config目录下的server文件和zookeeper文件,其中分别修改kafka的日志保存路径和zookeeper的数据保存路径。
在这里插入图片描述

启动kafka

先启动kafka自带的zookeeper,在kafka的根目录下打开终端,使用配置文件启动

./bin/windows/zookeeper-server-start.bat config/zookeeper.properties

同样在kafka目录的根目录下启动kafka

./bin/windows/kafka-server-start.bat config/server.properties

2 使用golang的github.com/Shopify/sarama库连接kafka

package main

import (
	"fmt"
	"time"

	"github.com/Shopify/sarama"
)

func main() {
	config:=sarama.NewConfig()
	// 生产者配置
	config.Producer.RequiredAcks=sarama.WaitForAll
	config.Producer.Partitioner=sarama.NewRandomPartitioner
	config.Producer.Return.Successes=true
	// 封装消息
	msg:=&sarama.ProducerMessage{}
	msg.Topic="shopping"
	time_str:=time.Now().Format("2006-01-02 15:04:05")
	msg.Value=sarama.StringEncoder("0413 test log!"+time_str)
	// 连接kafka
	client,err:=sarama.NewSyncProducer([]string{"127.0.0.1:9092"}, config)
	if err!=nil {
		fmt.Println("producer closed", err)
		return
	}
	defer client.Close()
	// 发送消息
	partition,offset,err:=client.SendMessage(msg)
	if err!=nil {
		fmt.Println("send failed", err)
		return
	}
	fmt.Printf("partition:%v offset:%v", partition, offset)
}

这段代码实现了模拟生产者向kafka发送消息的过程,包含:配置生产者,封装消息,消息类型是 *sarama.ProducerMessage,连接kafka,默认端口是9092,发送消息,返回消息存储的partition和offset日志偏移量。

3 确认生产者发送成功

使用kafka自带的命令行消费者客户端查看kafka中的数据
在kafka的根目录下

bin/windows/kafka-console-consumer.bat --bootstrap-server 127.0.0.1:9092 --topic shopping --from-beginning

这里的topic和代码中的topic一致,均为shopping
终端会输出之前发送的数据。

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 举报,一经查实,本站将立刻删除。

文章由极客之音整理,本文链接:https://www.bmabk.com/index.php/post/133471.html

(0)
飞熊的头像飞熊bm

相关推荐

发表回复

登录后才能评论
极客之音——专业性很强的中文编程技术网站,欢迎收藏到浏览器,订阅我们!