增加发送消息的功能
This commit is contained in:
@@ -5,6 +5,7 @@ import (
|
||||
"context"
|
||||
"flag"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/signal"
|
||||
@@ -337,6 +338,41 @@ func destroy() {
|
||||
log.Warn(" terminated ")
|
||||
}
|
||||
|
||||
func produceHandler(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
topic := vars["topic"]
|
||||
key := vars["key"]
|
||||
msg, err := ioutil.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
log.Error(err)
|
||||
http.Error(w, err.Error(), 500)
|
||||
return
|
||||
}
|
||||
writer := kafka.NewWriter(kafka.WriterConfig{
|
||||
Brokers: strings.Split(*bootstrapServers, ","),
|
||||
Topic: topic,
|
||||
Balancer: &kafka.LeastBytes{},
|
||||
})
|
||||
log.WithFields(log.Fields{
|
||||
"topic": topic,
|
||||
"key": key,
|
||||
"msg": string(msg),
|
||||
}).Info("sending message")
|
||||
senderr := writer.WriteMessages(context.Background(),
|
||||
kafka.Message{
|
||||
Key: []byte(key),
|
||||
Value: msg,
|
||||
})
|
||||
if senderr != nil {
|
||||
log.Error(err)
|
||||
http.Error(w, err.Error(), 500)
|
||||
return
|
||||
}
|
||||
log.Info("send ok")
|
||||
writer.Close()
|
||||
w.WriteHeader(200)
|
||||
}
|
||||
|
||||
func main() {
|
||||
flag.Parse()
|
||||
if *logfile != "" {
|
||||
@@ -360,6 +396,8 @@ func main() {
|
||||
|
||||
router.Path("/ws/{topic}/{partition}").Queries("offset", "{offset}").HandlerFunc(websocketHandler).Name("web-socket")
|
||||
|
||||
router.Path("/produce/{topic}/{key}").HandlerFunc(produceHandler).Methods("POST")
|
||||
|
||||
// Run our server in a goroutine so that it doesn't block.
|
||||
go func() {
|
||||
log.WithFields(log.Fields{"address": *addr}).Info("starting server")
|
||||
|
||||
Reference in New Issue
Block a user