NOTE

package

1. 是什么 go组织源码的方式 2. 包初始化流程 package2 test.go service.go dao.go 输出 3. 循环依赖 package1 test.go dao.go service.go 输出 可以看出 test/package1/dao 导入了 test/package1/service ,而 test/package1/service 导入了 test/package1/dao ,因此循环依赖了 4. 常

Go约 2 分钟读完historical

这是历史学习笔记,可能存在过时或不完整的理解。

1. 是什么

go组织源码的方式

2. 包初始化流程

  • package2_test.go
package package2

import (
	_ "test/package2/service"
	"testing"
)

func TestPackage1(t *testing.T) {

}
  • service.go
package service

import (
	"fmt"
	_ "test/package2/dao"
)

const ServiceName = "UserService"

var Type = "service"

func init() {
	fmt.Println(ServiceName, Type)
}
  • dao.go
package dao

import (
	"fmt"
)

const DaoName = "UserDao"

var Type = "dao"

func init() {
	fmt.Println(DaoName, Type)
}
  • 输出
UserDao dao
UserService service

3. 循环依赖

  • package1_test.go
package package1

import (
	_ "test/package1/dao"
	_ "test/package1/service"

	"testing"
)

func TestPackage1(t *testing.T) {

}
  • dao.go
package dao

import (
	"fmt"
	"test/package1/service"
)

const DaoName = "dao"

type Dao struct {
}

func (d Dao) GetUser() {
	fmt.Println(service.ServiceName)
}
  • service.go
package service

import (
	"fmt"
	"test/package1/dao"
)

const ServiceName = "service"

type Service struct {
}

func (s Service) GetUser() {
	fmt.Println(dao.DaoName)
}
  • 输出
import cycle not allowed
package test/package1 (test)
	imports test/package1/dao
	imports test/package1/service
	imports test/package1/dao

可以看出test/package1/dao导入了test/package1/service,而test/package1/service导入了test/package1/dao,因此循环依赖了

4. 常用的包

4.1. context

4.2. http

4.3. 数据库

4.4. 序列化

4.4.1. JSON

tidwall/gjson: Get JSON values quickly - JSON parser for Go 深入 Go 中各个高性能 JSON 解析库 - luozhiyun`s Blog

4.5. I/O

4.6. log

4.7. redis

4.8. kafka

4.8.1. sarama

  • Shopify/sarama: Sarama is a Go library for Apache Kafka 0.8, and up.

  • go操作kafka | 李文周的博客

  • 例子

    • producer
    const TopicName = "TopicName"
    
    func main() {
    	config := sarama.NewConfig()
    	config.Producer.RequiredAcks = sarama.WaitForAll          // 发送完数据需要leader和follow都确认
    	config.Producer.Partitioner = sarama.NewRandomPartitioner // 新选出一个partition
    	config.Producer.Return.Successes = true                   // 成功交付的消息将在success channel返回
    
    	// 连接kafka
    	producer, err := sarama.NewSyncProducer([]string{"127.0.0.1:9093"}, config)
    	if err != nil {
    		fmt.Println("producer closed, err:", err)
    		return
    	}
    	defer producer.Close()
    
    	i := 0
    	for {
    
    		// 构造一个消息
    		msg := &sarama.ProducerMessage{}
    		msg.Topic = TopicName
    		msg.Value = sarama.StringEncoder(fmt.Sprintf("this is a test log %d", i))
    		// 发送消息
    		pid, offset, err := producer.SendMessage(msg)
    		if err != nil {
    			fmt.Println("send msg failed, err:", err)
    			return
    		}
    		fmt.Printf("pid:%v offset:%v\n", pid, offset)
    		i++
    		time.Sleep(time.Second)
    	}
    
    }
    • consumer
    
    const TopicName = "TopicName"
    
    func main() {
    	// When
    	consumer, err := sarama.NewConsumer([]string{"localhost:9093"}, nil)
    	if err != nil {
    		fmt.Println(err)
    		return
    	}
    	defer consumer.Close()
    
    	partitionList, err := consumer.Partitions(TopicName) // 根据topic取到所有的分区
    	if err != nil {
    		fmt.Println(err)
    		return
    	}
    	fmt.Println(partitionList)
    
    	for _, partition := range partitionList { // 遍历所有的分区
    
    		go func(partition int32) {
    
    			consumePartition, err := consumer.ConsumePartition(TopicName, partition,
    				sarama.OffsetNewest)
    			if err != nil {
    				fmt.Println(err)
    			}
    			defer consumePartition.Close()
    			for {
    				select {
    				case message := <-consumePartition.Messages():
    					fmt.Printf("topic: %v, partition: %v, offset: %v, key: %v, value: %v\n",
    						message.Topic, message.Partition,
    						message.Offset, message.Key,
    						string(message.Value))
    				case err := <-consumePartition.Errors():
    					fmt.Println(err)
    				}
    			}
    
    		}(partition)
    
    	}
    
    	select {}
    }

4.8.2. confluent-kafka-go

Home · edenhill/librdkafka Wiki · GitHub Kafka Go Client | Confluent Documentation kafka - Go Documentation Server Consumer Configurations | Confluent Documentation

func replayKafkaEvent(ctx context.Context, totalInvites int64, rankID string) {
	c, err := kafka.NewConsumer(&kafka.ConfigMap{
		"bootstrap.servers":                  "<KAFKA_BROKER>:19092",
		"group.id":                           uuid.New().String(),
		"auto.offset.reset":                  "earliest",
		"topic.metadata.refresh.interval.ms": 5,
	})

	if err != nil {
		log.FatalContext(ctx, err)
	}
	topic := fmt.Sprintf("%s%v", "update_cache_", rankID)
	err = c.SubscribeTopics([]string{topic}, nil)
	if err != nil {
		log.FatalContext(ctx, err)
	}

	messages := make([]*kafka.Message, 0, 100000)
	for {
		msg, err := c.ReadMessage(-1)
		if err == nil {
			messages = append(messages, msg)
			chs <- msg
			//if len(messages) == 100000 {
			//	invites, _ := BatchUpdateCacheMsgConsumer(ctx, messages)
			//	if totalInvites <= invites {
			//		log.InfoContext(ctx, "rankID:%v init ok, invites:%v", rankID, invites)
			//		once.Do(func() {
			//			InitHotRankCacheFinished <- 1
			//		})
			//	}
			//	messages = make([]*kafka.Message, 0, 100000)
			//}
		} else {
			log.ErrorContextf(ctx, "Consumer error: %v (%v)\n", err, msg)
		}
	}
	//go func(c *kafka.Consumer) {
	//	partitions, err := c.Assignment()
	//	if err != nil {
	//		fmt.Fprintf(os.Stderr, "设置topic失败: %s\n", err)
	//		os.Exit(1)
	//	}
	//	partitions, err = c.OffsetsForTimes(partitions, ts.Second()*100000)
	//	if err != nil {
	//		fmt.Fprintf(os.Stderr, "设置topic失败: %s\n", err)
	//		os.Exit(1)
	//	}
	//	for _, partition := range partitions {
	//		c.Seek(partition, 0)
	//	}
	//	for {
	//		msg, err := c.ReadMessage(-1)
	//		if err == nil {
	//			messages = append(messages, msg)
	//			if len(messages) == 100000 {
	//				BatchUpdateCacheMsgConsumer(ctx, messages)
	//				messages = make([]*kafka.Message, 0, 100000)
	//			}
	//			ts = msg.Timestamp
	//			if totalInvites == 100000 {
	//				break
	//			}
	//		} else {
	//			log.Errorf("Consumer error: %v (%v)\n", err, msg)
	//		}
	//	}
	//	c.Close()
	//}(c)

}

4.9. cache

4.10. 格式化

  • 安装
go get -u github.com/cuonglm/gocmt
  • 在项目根目录下执行
gocmt -d ./ -i

4.11. 协程池

alitto/pond: 🔘 Minimalistic and High-performance goroutine worker pool written in Go

4.12. error

golang/xerrors xerrors package - golang.org/x/xerrors - Go Packages

4.13. errgroup

errgroup package - golang.org/x/sync/errgroup - Go Packages sync/errgroup at master · golang/sync golang 详解协程——errgroup - CodeAntenna 聊聊 ErrorGroup 的用法和拓展 | Mark’s Blog

4.14. awesome go

Awesome Go | LibHunt

5. 参考

讨论

使用 GitHub 账号参与讨论,评论会保存在 GitHub Issues 中。在 GitHub 查看