diff --git a/docs/develop-with-kafka-api/_index.md b/docs/develop-with-kafka-api/_index.md index 396af55..161a411 100644 --- a/docs/develop-with-kafka-api/_index.md +++ b/docs/develop-with-kafka-api/_index.md @@ -1,5 +1,5 @@ --- -order: ['compatibility.md', 'java.md', 'python.md'] +order: ['compatibility.md', 'java.md', 'python.md', 'go.md'] collapsed: false --- diff --git a/docs/develop-with-kafka-api/compatibility.md b/docs/develop-with-kafka-api/compatibility.md index 5ecc334..fd02666 100644 --- a/docs/develop-with-kafka-api/compatibility.md +++ b/docs/develop-with-kafka-api/compatibility.md @@ -27,7 +27,7 @@ Currenty, the clients below are tested by HStream. | -------- | ----------------------------------------------------------- | | Java | [Apache Kafka Java Client](https://github.com/apache/kafka) | | Python | [kafka-python], [confluent-kafka-python] | -| Go | [franz-go](https://github.com/twmb/franz-go) | +| Go | [kafka-go](https://github.com/segmentio/kafka-go) | | C/C++ | [librdkafka](https://github.com/confluentinc/librdkafka) | ::: tip diff --git a/docs/develop-with-kafka-api/go.md b/docs/develop-with-kafka-api/go.md new file mode 100644 index 0000000..8333433 --- /dev/null +++ b/docs/develop-with-kafka-api/go.md @@ -0,0 +1,31 @@ +# Develop with Go Kafka client + +This page shows how to use [Kafka-go Client](https://github.com/segmentio/kafka-go) to interact with HStream. + + +## Create a Topic + + +::: code-group + +<<< @/../kafka-examples/go/examples/create_topics.go [Go] + +::: + +## Produce Records + + +::: code-group + +<<< @/../kafka-examples/go/examples/produce.go [Go] + +::: + +## Consume Records + + +::: code-group + +<<< @/../kafka-examples/go/examples/consume.go [Go] + +::: diff --git a/kafka-examples/go/examples/common.go b/kafka-examples/go/examples/common.go new file mode 100644 index 0000000..6dea8eb --- /dev/null +++ b/kafka-examples/go/examples/common.go @@ -0,0 +1,5 @@ +package examples + +var ( + totalMesssages = 10 +) diff --git a/kafka-examples/go/examples/consume.go b/kafka-examples/go/examples/consume.go new file mode 100644 index 0000000..bf66481 --- /dev/null +++ b/kafka-examples/go/examples/consume.go @@ -0,0 +1,38 @@ +package examples + +import ( + "context" + "log" + "time" + + "github.com/segmentio/kafka-go" +) + +func Consume() { + host := "localhost:9092" + + reader := kafka.NewReader(kafka.ReaderConfig{ + Brokers: []string{host}, + Topic: "test-topic", + GroupID: "test-group1", + StartOffset: kafka.FirstOffset, + }) + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + readCnt := 0 + for readCnt < totalMesssages { + m, err := reader.ReadMessage(ctx) + if err != nil { + log.Fatal(err) + } + readCnt++ + log.Printf("Message received: value = %s, timestamp = %v, topic = %s", string(m.Value), m.Time, m.Topic) + } + log.Printf("Read %d messages", readCnt) + + if err := reader.Close(); err != nil { + log.Fatal("Failed to close reader:", err) + } +} diff --git a/kafka-examples/go/examples/create_topics.go b/kafka-examples/go/examples/create_topics.go new file mode 100644 index 0000000..a06099a --- /dev/null +++ b/kafka-examples/go/examples/create_topics.go @@ -0,0 +1,40 @@ +package examples + +import ( + "context" + "errors" + "log" + "time" + + "github.com/segmentio/kafka-go" +) + +func CreateTopics() { + host := "localhost:9092" + client := &kafka.Client{ + Addr: kafka.TCP(host), + Timeout: 10 * time.Second, + } + + request := &kafka.CreateTopicsRequest{ + Topics: []kafka.TopicConfig{ + { + Topic: "test-topic", + NumPartitions: 1, + ReplicationFactor: 1, + }, + }, + ValidateOnly: false, + } + + resp, err := client.CreateTopics(context.Background(), request) + if err != nil { + log.Fatal(err) + } + + for _, err = range resp.Errors { + if err != nil && !errors.Is(err, kafka.TopicAlreadyExists) { + log.Fatal(err) + } + } +} diff --git a/kafka-examples/go/examples/produce.go b/kafka-examples/go/examples/produce.go new file mode 100644 index 0000000..637a05a --- /dev/null +++ b/kafka-examples/go/examples/produce.go @@ -0,0 +1,55 @@ +package examples + +import ( + "context" + "fmt" + "log" + "sync" + + "github.com/segmentio/kafka-go" +) + +func Produce() { + host := "localhost:9092" + + wg := &sync.WaitGroup{} + wg.Add(totalMesssages) + writer := &kafka.Writer{ + Addr: kafka.TCP(host), + Topic: "test-topic", + Balancer: &kafka.RoundRobin{}, + RequiredAcks: kafka.RequireAll, + Async: true, + Completion: func(messages []kafka.Message, err error) { + if err != nil { + wg.Done() + log.Printf("produce err: %s\n", err.Error()) + return + } + + for _, msg := range messages { + wg.Done() + log.Printf("write date to partition %d, offset %d\n", msg.Partition, msg.Offset) + } + }, + } + + defer func() { + if err := writer.Close(); err != nil { + log.Fatal("Failed to close writer:", err) + } + }() + + for i := 0; i < totalMesssages; i++ { + msg := kafka.Message{ + Key: []byte(fmt.Sprintf("key-%d", i)), + Value: []byte(fmt.Sprintf("value-%d", i)), + } + if err := writer.WriteMessages(context.Background(), msg); err != nil { + log.Fatal("Failed to write messages:", err) + } + } + + wg.Wait() + log.Println("Write messages done.") +} diff --git a/kafka-examples/go/go.mod b/kafka-examples/go/go.mod new file mode 100644 index 0000000..b256e31 --- /dev/null +++ b/kafka-examples/go/go.mod @@ -0,0 +1,11 @@ +module github.com/hstreamdb/hstream-kafka-go-examples + +go 1.21 + +require github.com/segmentio/kafka-go v0.4.47 + +require ( + github.com/klauspost/compress v1.16.7 // indirect + github.com/pierrec/lz4/v4 v4.1.18 // indirect + github.com/stretchr/testify v1.8.4 // indirect +) diff --git a/kafka-examples/go/go.sum b/kafka-examples/go/go.sum new file mode 100644 index 0000000..b62c751 --- /dev/null +++ b/kafka-examples/go/go.sum @@ -0,0 +1,71 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= +github.com/klauspost/compress v1.16.7 h1:2mk3MPGNzKyxErAw8YaohYh69+pa4sIQSC0fPGCFR9I= +github.com/klauspost/compress v1.16.7/go.mod h1:ntbaceVETuRiXiv4DpjP66DpAtAGkEQskQzEyD//IeE= +github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= +github.com/pierrec/lz4/v4 v4.1.18 h1:xaKrnTkyoqfh1YItXl56+6KJNVYWlEEPuAQW9xsplYQ= +github.com/pierrec/lz4/v4 v4.1.18/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/segmentio/kafka-go v0.4.47 h1:IqziR4pA3vrZq7YdRxaT3w1/5fvIH5qpCwstUanQQB0= +github.com/segmentio/kafka-go v0.4.47/go.mod h1:HjF6XbOKh0Pjlkr5GVZxt6CsjjwnmhVOfURM5KMd8qg= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= +github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.4 h1:CcVxjf3Q8PM0mHUKJCdn+eZZtm5yQwehR5yeSVQQcUk= +github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= +github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c= +github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= +github.com/xdg-go/scram v1.1.2 h1:FHX5I5B4i4hKRVRBCFRxq1iQRej7WO3hhBuJf+UUySY= +github.com/xdg-go/scram v1.1.2/go.mod h1:RT/sEzTbU5y00aCK8UOx6R7YryM0iF1N2MOmC3kKLN4= +github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8= +github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM= +github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= +golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= +golang.org/x/crypto v0.14.0/go.mod h1:MVFd36DqK4CsrnJYDkBA3VC4m2GkXAM0PvzMCn4JQf4= +golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= +golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= +golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= +golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= +golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= +golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg= +golang.org/x/net v0.17.0 h1:pVaXccu2ozPjCXewfr1S7xza/zcXTity9cCdXQYSjIM= +golang.org/x/net v0.17.0/go.mod h1:NxSsAGuq816PNPmqtQdLE42eU2Fs7NoRIZrHJAlaCOE= +golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= +golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= +golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= +golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= +golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo= +golang.org/x/term v0.13.0/go.mod h1:LTmsnFJwVN6bCy1rVCoS+qHT1HhALEFxKncY3WNNh4U= +golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= +golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= +golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= +golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ= +golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= +golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8= +golang.org/x/text v0.13.0 h1:ablQoSUd0tRdKxZewP80B+BaqeKJuVhuRxj/dkrun3k= +golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE= +golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= +golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= +golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= +golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/kafka-examples/go/main.go b/kafka-examples/go/main.go new file mode 100644 index 0000000..0d8d97d --- /dev/null +++ b/kafka-examples/go/main.go @@ -0,0 +1,33 @@ +package main + +import ( + "log" + "reflect" + "runtime" + + "github.com/hstreamdb/hstream-kafka-go-examples/examples" +) + +func main() { + xs := []func(){ + examples.CreateTopics, + examples.Produce, + examples.Consume, + } + + for _, x := range xs { + runFuncWithLog(x) + } + +} + +func getFuncName(i interface{}) string { + return runtime.FuncForPC(reflect.ValueOf(i).Pointer()).Name() +} + +func runFuncWithLog(f func()) { + funcName := getFuncName(f) + log.Printf("==== start %s ====", funcName) + f() + log.Printf("==== end %s ====", funcName) +}