1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95
|
package main
import ( "context" "fmt" "github.com/segmentio/kafka-go" "os" "os/signal" "syscall" "time" )
var topic = "test-topic" var reader *kafka.Reader
func writeKafka(ctx context.Context) { writer := &kafka.Writer{ Addr: kafka.TCP("localhost:9092"), Topic: topic, Balancer: &kafka.Hash{}, WriteTimeout: time.Second * 1, RequiredAcks: kafka.RequireNone, AllowAutoTopicCreation: true, } defer writer.Close() for i := 1; i <= 5; i++ { if err := writer.WriteMessages(ctx, kafka.Message{Key: []byte("Key-A"), Value: []byte("this")}, kafka.Message{Key: []byte("Key-B"), Value: []byte("is")}, kafka.Message{Key: []byte("Key-C"), Value: []byte("a")}, kafka.Message{Key: []byte("Key-D"), Value: []byte("test")}, ); err != nil { if err == kafka.LeaderNotAvailable { time.Sleep(time.Millisecond * 500) continue } else { fmt.Printf("failed to write message:%v\n", err) } } else { fmt.Printf("kafka write message success:%d\n", i) break } } }
func readKafka(ctx context.Context) { reader = kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"localhost:9092"}, Topic: topic, CommitInterval: time.Second * 1, GroupID: "consumer-group-id", StartOffset: kafka.FirstOffset, }) defer reader.Close() for { if message, err := reader.ReadMessage(ctx); err != nil { fmt.Printf("failed to read message:%v\n", err) break } else { fmt.Printf("message at topic:%s,partition:%d,offset:%d,key:%s,value:%s\n", message.Topic, message.Partition, message.Offset, string(message.Key), string(message.Value)) } } }
func listenSignal() { quit := make(chan os.Signal, 1) signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) sig := <-quit fmt.Printf("recv:%s\n", sig.String()) if reader != nil { reader.Close() } os.Exit(0) }
func main() { ctx := context.Background() writeKafka(ctx) go listenSignal() readKafka(ctx) }
|