NOTE

package

1. What it is How Go organizes source code 2. Package initialization flow 3. Circular dependencies 4. Common packages 5. References

Go1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What It Is

The way Go organizes source code.

2. Package Initialization Flow

  • 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)
}
  • Output
UserDao dao
UserService service

3. Circular Dependencies

  • 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)
}
  • Output
import cycle not allowed
package test/package1 (test)
	imports test/package1/dao
	imports test/package1/service
	imports test/package1/dao

As you can see, test/package1/dao imports test/package1/service, while test/package1/service imports test/package1/dao, which creates a circular dependency.

4. Common Packages

4.1. context

4.2. http

4.3. Database

4.4. Serialization

4.4.1. JSON

tidwall/gjson: Get JSON values quickly - JSON parser for Go Deep Dive into High-Performance JSON Libraries in Go - 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.
  • Using Kafka with Go | Li Wenzhou’s Blog
  • Examples
    • producer

      const TopicName = "TopicName"
      
      func main() {
      	config := sarama.NewConfig()
      	config.Producer.RequiredAcks = sarama.WaitForAll          // leader and followers must all acknowledge after data is sent
      	config.Producer.Partitioner = sarama.NewRandomPartitioner // choose a new partition
      	config.Producer.Return.Successes = true                   // successfully delivered messages are returned on the success channel
      
      	// Connect to 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 {
      		// Construct a message.
      		msg := &sarama.ProducerMessage{}
      		msg.Topic = TopicName
      		msg.Value = sarama.StringEncoder(fmt.Sprintf("this is a test log %d", i))
      		// Send the message.
      		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) // get all partitions for the topic
      	if err != nil {
      		fmt.Println(err)
      		return
      	}
      	fmt.Println(partitionList)
      
      	for _, partition := range partitionList { // iterate over all partitions
      		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, "failed to set topic: %s\n", err)
	//		os.Exit(1)
	//	}
	//	partitions, err = c.OffsetsForTimes(partitions, ts.Second()*100000)
	//	if err != nil {
	//		fmt.Fprintf(os.Stderr, "failed to set 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. Formatting

  • Installation
go get -u github.com/cuonglm/gocmt
  • Run from the project root directory
gocmt -d ./ -i

4.11. Goroutine Pool

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 Goroutines Explained: errgroup - CodeAntenna ErrorGroup Usage and Extensions | Mark’s Blog

4.14. awesome go

Awesome Go | LibHunt

5. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub