Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 47 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
module github.com/pavedroad-io/go-core

go 1.17

require (
github.com/Shopify/sarama v1.32.0
github.com/gofrs/uuid v4.2.0+incompatible
github.com/sirupsen/logrus v1.8.1
github.com/spf13/viper v1.10.1
go.uber.org/zap v1.21.0
gopkg.in/natefinch/lumberjack.v2 v2.0.0
gopkg.in/yaml.v2 v2.4.0
)

require (
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/eapache/go-resiliency v1.2.0 // indirect
github.com/eapache/go-xerial-snappy v0.0.0-20180814174437-776d5712da21 // indirect
github.com/eapache/queue v1.1.0 // indirect
github.com/fsnotify/fsnotify v1.5.1 // indirect
github.com/golang/snappy v0.0.4 // indirect
github.com/hashicorp/go-uuid v1.0.2 // indirect
github.com/hashicorp/hcl v1.0.0 // indirect
github.com/jcmturner/aescts/v2 v2.0.0 // indirect
github.com/jcmturner/dnsutils/v2 v2.0.0 // indirect
github.com/jcmturner/gofork v1.0.0 // indirect
github.com/jcmturner/gokrb5/v8 v8.4.2 // indirect
github.com/jcmturner/rpc/v2 v2.0.3 // indirect
github.com/klauspost/compress v1.14.4 // indirect
github.com/magiconair/properties v1.8.5 // indirect
github.com/mitchellh/mapstructure v1.4.3 // indirect
github.com/pelletier/go-toml v1.9.4 // indirect
github.com/pierrec/lz4 v2.6.1+incompatible // indirect
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 // indirect
github.com/spf13/afero v1.6.0 // indirect
github.com/spf13/cast v1.4.1 // indirect
github.com/spf13/jwalterweatherman v1.1.0 // indirect
github.com/spf13/pflag v1.0.5 // indirect
github.com/subosito/gotenv v1.2.0 // indirect
go.uber.org/atomic v1.7.0 // indirect
go.uber.org/multierr v1.6.0 // indirect
golang.org/x/crypto v0.0.0-20220214200702-86341886e292 // indirect
golang.org/x/net v0.0.0-20220127200216-cd36cc0744dd // indirect
golang.org/x/sys v0.0.0-20211216021012-1d35b9e2eb4e // indirect
golang.org/x/text v0.3.7 // indirect
gopkg.in/ini.v1 v1.66.2 // indirect
)
844 changes: 844 additions & 0 deletions go.sum

Large diffs are not rendered by default.

83 changes: 83 additions & 0 deletions logger/examples/general-consumer/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
package main

import (
"fmt"
"os"
"os/signal"
"sync"

"github.com/Shopify/sarama"
)

var (
brokers = []string{"localhost:9092"}
config = sarama.NewConfig()
)

func main() {

// init master consumer
master, err := sarama.NewConsumer(brokers, config)
if err != nil {
fmt.Fprintf(os.Stderr, "Failed to get consumer: %s\n", err.Error())
os.Exit(1)
}
defer master.Close()
// get topics
topics, err := master.Topics()
if err != nil {
fmt.Fprintf(os.Stderr, "Failed to get topics: %s\n", err.Error())
os.Exit(1)
}
// init exit
interrupt := make(chan os.Signal, 1)
signal.Notify(interrupt, os.Interrupt)
wg := sync.WaitGroup{}

// consume all topics
for _, topic := range topics {
// get partitions
partitions, err := master.Partitions(topic)
if err != nil {
fmt.Fprintf(os.Stderr,
"Failed to get partitions for topic %s: %s\n",
topic, err.Error())
os.Exit(1)
}
// consume all partitions
for _, partition := range partitions {
// init partition consumer
consumer, err := master.ConsumePartition(topic, partition,
sarama.OffsetOldest)
if err != nil {
fmt.Fprintf(os.Stderr,
"Failed to get consumer for topic %s partition %+v: %s\n",
topic, partition, err.Error())
os.Exit(1)
}
defer consumer.Close()
// consume log messages
wg.Add(1)
go func(w *sync.WaitGroup) {
loop:
for {
select {
case msg := <-consumer.Messages():
fmt.Printf("T:%s P:%d K:%s V:%s\n\n",
msg.Topic, msg.Partition, msg.Key, msg.Value)
case err := <-consumer.Errors():
fmt.Fprintf(os.Stderr, "Error T:%s P:%+v %s\n",
topic, partition, err)
break loop
case <-interrupt:
fmt.Printf("\nExiting consumer\n")
break loop
}
}
w.Done()
return
}(&wg)
}
}
wg.Wait()
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,42 +7,44 @@ import (
"os/signal"

"github.com/Shopify/sarama"
cluster "github.com/bsm/sarama-cluster"
)

var (
brokers = []string{"localhost:9092"}
group = "mygroup"
topics = []string{"logs"}
config = cluster.NewConfig()
brokers = []string{"localhost:9092"}
topic string = "logs"
config = sarama.NewConfig()
)

// Use sarama-cluster to consume logs

func main() {

// init config
config.Consumer.Offsets.Initial = sarama.OffsetOldest

// init consumer
consumer, err := cluster.NewConsumer(brokers, group, topics, config)
master, err := sarama.NewConsumer(brokers, config)
if err != nil {
panic(err)
fmt.Fprintf(os.Stderr, "Failed to get consumer: %s\n", err.Error())
os.Exit(1)
}
// init partition 0 consumer
consumer, err := master.ConsumePartition(topic, 0, sarama.OffsetOldest)
if err != nil {
fmt.Fprintf(os.Stderr, "Failed to get partition consumer: %s\n",
err.Error())
os.Exit(1)
}
// uncomment to enable connection debugging
// sarama.Logger = stdlog.New(os.Stdout, "[sarama] ", stdlog.LstdFlags)

// init exit
interrupt := make(chan os.Signal, 1)
signal.Notify(interrupt, os.Interrupt)

defer consumer.Close()
defer master.Close()

// consume log messages
for {
select {
case msg := <-consumer.Messages():
fmt.Printf("P:%d K:%s V:%s\n\n", msg.Partition, msg.Key, msg.Value)
consumer.MarkOffset(msg, "kafka-test")
case err := <-consumer.Errors():
fmt.Fprintf(os.Stderr, "Message error: %s\n", err)
case <-interrupt:
Expand Down
14 changes: 13 additions & 1 deletion logger/kafka.go
Original file line number Diff line number Diff line change
Expand Up @@ -285,9 +285,21 @@ func (kp *KafkaProducer) sendMessage(msg []byte) error {
}

kp.producer.Input() <- &sarama.ProducerMessage{
Key: key,
Topic: topic.(string),
Key: key,
Value: sarama.ByteEncoder(newmsg),
}
return nil
}

// sendMessageTK send msg with topic and key and no processing
func (kp *KafkaProducer) sendMessageTKV(topic string, key []byte,
value []byte) error {

kp.producer.Input() <- &sarama.ProducerMessage{
Topic: topic,
Key: sarama.ByteEncoder(key),
Value: sarama.ByteEncoder(value),
}
return nil
}
Loading