Skip to content

Commit

Permalink
Configurable kafka protocol/server version in producer
Browse files Browse the repository at this point in the history
Signed-off-by: Chodor Marek <[email protected]>
  • Loading branch information
Chodor Marek committed Jul 9, 2019
1 parent 66ccdac commit 9a90ad7
Show file tree
Hide file tree
Showing 2 changed files with 25 additions and 11 deletions.
10 changes: 9 additions & 1 deletion pkg/kafka/producer/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@ type Builder interface {

// Configuration describes the configuration properties needed to create a Kafka producer
type Configuration struct {
Brokers []string
Brokers []string
ProtocolVersion string
auth.AuthenticationConfig
}

Expand All @@ -36,5 +37,12 @@ func (c *Configuration) NewProducer() (sarama.AsyncProducer, error) {
saramaConfig := sarama.NewConfig()
saramaConfig.Producer.Return.Successes = true
c.AuthenticationConfig.SetConfiguration(saramaConfig)
if len(c.ProtocolVersion) > 0 {
ver, err := sarama.ParseKafkaVersion(c.ProtocolVersion)
if err != nil {
return nil, err
}
saramaConfig.Version = ver
}
return sarama.NewAsyncProducer(c.Brokers, saramaConfig)
}
26 changes: 16 additions & 10 deletions plugin/storage/kafka/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,14 @@ const (
// EncodingZipkinThrift is used for spans encoded as Zipkin Thrift.
EncodingZipkinThrift = "zipkin-thrift"

configPrefix = "kafka.producer"
suffixBrokers = ".brokers"
suffixTopic = ".topic"
suffixEncoding = ".encoding"
defaultBroker = "127.0.0.1:9092"
defaultTopic = "jaeger-spans"
defaultEncoding = EncodingProto
configPrefix = "kafka.producer"
suffixBrokers = ".brokers"
suffixTopic = ".topic"
suffixProtocolVersion = ".protocol-version"
suffixEncoding = ".encoding"
defaultBroker = "127.0.0.1:9092"
defaultTopic = "jaeger-spans"
defaultEncoding = EncodingProto
)

var (
Expand All @@ -49,9 +50,9 @@ var (

// Options stores the configuration options for Kafka
type Options struct {
config producer.Configuration
topic string
encoding string
config producer.Configuration
topic string
encoding string
}

// AddFlags adds flags for Options
Expand All @@ -64,6 +65,10 @@ func (opt *Options) AddFlags(flagSet *flag.FlagSet) {
configPrefix+suffixTopic,
defaultTopic,
"(experimental) The name of the kafka topic")
flagSet.String(
configPrefix+suffixProtocolVersion,
"",
"(experimental) Kafka protocol version - must be supported by kafka server")
flagSet.String(
configPrefix+suffixEncoding,
defaultEncoding,
Expand All @@ -78,6 +83,7 @@ func (opt *Options) InitFromViper(v *viper.Viper) {
authenticationOptions.InitFromViper(configPrefix, v)
opt.config = producer.Configuration{
Brokers: strings.Split(stripWhiteSpace(v.GetString(configPrefix+suffixBrokers)), ","),
ProtocolVersion: v.GetString(configPrefix + suffixProtocolVersion),
AuthenticationConfig: authenticationOptions,
}
opt.topic = v.GetString(configPrefix + suffixTopic)
Expand Down

0 comments on commit 9a90ad7

Please sign in to comment.