brokers := []string{"xxx"}
config := sarama.NewConfig()
config.Version = sarama.V0_10_2_0
config.Net.SASL.Enable = true
config.Net.SASL.User = "用户"
config.Net.SASL.Password = "密码"
config.Net.SASL.Mechanism = sarama.SASLTypePlaintext
config.Net.TLS.Enable = true
config.Net.SASL.Version = 0
config.Net.TLS.Config = &tls.Config{}
config.Net.TLS.Config.InsecureSkipVerify = true
var client sarama.Client
client, err = sarama.NewClient(brokers, config)
fmt.Println("create sarama client error:", err)
if err := client.Close(); err != nil {
fmt.Println("client close error:", err)
fmt.Println("Close client!")
consumer, err := sarama.NewConsumerFromClient(client)
if err := consumer.Close(); err != nil {
fmt.Println("consumer close error:", err)
fmt.Println("Close consumer!")
fmt.Println("new consumer error:", err)
offsetManager, err := sarama.NewOffsetManagerFromClient(cg, client)
if err := offsetManager.Close(); err != nil {
fmt.Println("close offset manager error:", err)
fmt.Println("Close offset manager!")
fmt.Println("new offset manager error:", err)
if !config.Consumer.Offsets.AutoCommit.Enable {
defer offsetManager.Commit()
signals := make(chan os.Signal, 1)
signal.Notify(signals, os.Interrupt)
partitions, err := client.Partitions(topic)
fmt.Println("get partitions error:", err)
fmt.Println("start consume. topic:", topic, "partitions:", partitions)
for _, partition := range partitions {
go consumePartition(quit, &thr, consumer, offsetManager, int32(partition), topic)
func consumePartition(quit <-chan int, thr *sync.WaitGroup, consumer sarama.Consumer,
offsetManager sarama.OffsetManager, partition int32, topic string) {
partitionOffsetManager, err := offsetManager.ManagePartition(topic, partition)
fmt.Println("create partitionOffsetManager err:", err)
offset, _ := partitionOffsetManager.NextOffset()
offset = sarama.OffsetNewest
var partitionConsumer sarama.PartitionConsumer
partitionConsumer, err = consumer.ConsumePartition(topic, partition, offset)
fmt.Println("create partition consumer err:", err)
if err := partitionConsumer.Close(); err != nil {
fmt.Println("Close partition", partition)
case msg := <-partitionConsumer.Messages():
fmt.Println("[Consume] topic:",
msg.Topic, "parition:", msg.Partition, "offset:", msg.Offset, "key:", msg.Key, "value:", msg.Value)
partitionOffsetManager.MarkOffset(msg.Offset+1, "")
case err := <-partitionConsumer.Errors():
fmt.Println("consume error:", err)