gpt4 book ai didi

docker - 当部署到 Docker 时,在 Golang 中实现的 Apache Kafka 消费者会出现 panic

转载 作者:数据小太阳 更新时间:2023-10-29 03:23:59 26 4
gpt4 key购买 nike

这是我尝试实现一个简单的微服务,它应该从 kafka 服务器读取消息并通过 HTTP 发送它。当我从终端运行它时它工作正常,但是当部署到 docker 上时它会出现

panic
panic: runtime error: invalid memory address or nil pointer dereference
[signal SIGSEGV: segmentation violation code=0x1 addr=0x40 pc=0x7b6345]

goroutine 12 [running]:
main.kafkaRoutine.func1(0xc420174060, 0x0, 0x0)
/go/src/github.com/deathcore666/ProperConsumerServiceYo/kafka.go:36 +0x95
created by main.kafkaRoutine
/go/src/github.com/deathcore666/ProperConsumerServiceYo/kafka.go:32 +0x1ad

kafka.go 第 32 和 36 行是 go func(pc sarama.PartitionConsumer) 函数所在的行。我对编程比较陌生,所以任何帮助将不胜感激。谢谢!

main.go:

func main() {
var (
listen = flag.String("listen", ":8080", "HTTP listen address")
proxy = flag.String("proxy", "", "Optional comma-separated list of URLs to proxy uppercase requests")
)
flag.Parse()

logger := log.NewLogfmtLogger(os.Stderr)

var svc KafkaService
svc = kafkaService{}
svc = proxyingMiddleware(context.Background(), *proxy, logger)(svc)
svc = loggingMiddleware(logger)(svc)


consumehandler := httptransport.NewServer(
makeConsumeEndpoint(svc),
decodeConsumeRequest,
encodeResponse,
)

http.Handle("/consume", consumehandler)

logger.Log("msg", "HTTP", "addr", *listen)
logger.Log("err", http.ListenAndServe(*listen, nil))}

服务.go:

    package main

import (
"context"
"errors"
"time"
)

//KafkaService yolo
type KafkaService interface {
Consume(context.Context, string) (string, error)
}

//ErrEmpty yolo
var ErrEmpty = errors.New("No topic provided")

type kafkaService struct{}

//Consumer logic implemented here
func (kafkaService) Consume(_ context.Context, topic string) (string, error) {
if topic == "" {
return "", ErrEmpty
}

var inChan = make(chan string)
var readyChan = make(chan struct{})
var result string
var brokers = []string{"192.168.88.208:9092"}
//var brokersLocal = []string{"localhost:9092"}
go kafkaRoutine(inChan, topic, brokers)
go func() {
for {
select {
case msg := <-inChan:
result = result + msg + "\n"
case <-time.After(time.Second * 1):
readyChan <- struct{}{}
}

}
}()

<-readyChan
close(inChan)
return result, nil
}

//ServiceMiddleware is a chainable thing for the service
type ServiceMiddleware func(KafkaService) KafkaService

kafka.go:

package main

import (
"fmt"
"time"

"github.com/Shopify/sarama"
)

func kafkaRoutine(inChan chan string, topic string, brokers []string) {
config := sarama.NewConfig()
config.Consumer.Return.Errors = true
consumer, err := sarama.NewConsumer(brokers, config)
if err != nil {
panic(err)
}

topics, _ := consumer.Topics()
if !(containsTopic(topics, topic)) {
inChan <- "There is no such a topic"
fmt.Println("kafkaroutine exited")
return
}

partitionList, err := consumer.Partitions(topic)
for _, partition := range partitionList {
pc, _ := consumer.ConsumePartition(topic, partition, sarama.OffsetOldest)
go func(pc sarama.PartitionConsumer) {
loop:
for {
select {
case msg := <-pc.Messages():
inChan <- string(msg.Value)
case <-time.After(time.Second * 1):
break loop
}
}
}(pc)
}
fmt.Println("Kafka GoRoutine exited")
}

func containsTopic(topics []string, topic string) bool {
for _, v := range topics {
if topic == v {
return true
}
}
return false
}

最佳答案

在 kafka.go 的第 27 行,您忽略了 ConsumePartition() 返回的错误。它很可能返回错误而不是有效的分区使用者,但由于您在尝试使用分区使用者时忽略了它,它会崩溃。

关于docker - 当部署到 Docker 时,在 Golang 中实现的 Apache Kafka 消费者会出现 panic ,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/47301748/

26 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com