featurebase/idk/kafka/cmd.go
Jacob Brinlee 1046e28338 SUP-288 (#2414)
* adding kafka consumer config options (--kafka-max-poll-interval, --kafka-session-timeout,  --kafka-group-instance-id, --kafka-socket-keepalive-enable, and --consumer-close-timeout)

* wrapping consumer.Close() in timeout. Will wait consumer-close-timeout seconds before forcing consumer to exit

* clean up logs

(cherry picked from commit 17cdc58d80)
2023-01-19 22:10:15 +00:00

63 lines
2.2 KiB
Go

package kafka
import (
"time"
"github.com/featurebasedb/featurebase/v3/idk"
"github.com/pkg/errors"
)
type Main struct {
idk.Main `flag:"!embed"`
idk.ConfluentCommand `flag:"!embed"`
Group string `help:"Kafka group."`
Topics []string `help:"Kafka topics to read from."`
Timeout time.Duration `help:"Time to wait for more records from Kafka before flushing a batch. 0 to disable."`
SkipOld bool `short:"" help:"False sets kafka consumer configuration auto.offset.reset to earliest, True sets it to latest."`
ConsumerCloseTimeout int `help:"The amount of time in seconds to wait for the consumer to close properly."`
}
func NewMain() (*Main, error) {
m := &Main{
Main: *idk.NewMain(),
ConfluentCommand: idk.ConfluentCommand{
KafkaBootstrapServers: []string{"localhost:9092"},
},
Group: "defaultgroup",
Topics: []string{"defaulttopic"},
Timeout: time.Second,
ConsumerCloseTimeout: 30,
}
m.SchemaRegistryURL = "http://" + defaultRegistryHost
m.Main.Namespace = "ingester_kafka"
//m.Main.OffsetMode = m.OffsetMode
m.OffsetMode = true
m.NewSource = func() (idk.Source, error) {
source := NewSource()
source.KafkaBootstrapServers = m.KafkaBootstrapServers
source.SchemaRegistryURL = m.SchemaRegistryURL
source.Group = m.Group
source.Topics = m.Topics
source.Log = m.Main.Log()
source.Timeout = m.Timeout
source.KafkaSocketTimeoutMs = int(m.Timeout / time.Millisecond)
source.SkipOld = m.SkipOld
source.ConfluentCommand = m.ConfluentCommand
source.SchemaRegistryUsername = m.SchemaRegistryUsername
source.SchemaRegistryPassword = m.SchemaRegistryPassword
source.Verbose = m.Verbose
source.KafkaMaxPollInterval = m.KafkaMaxPollInterval
source.KafkaSessionTimeout = m.KafkaSessionTimeout
source.KafkaGroupInstanceId = m.KafkaGroupInstanceId
source.KafkaDebug = m.KafkaDebug
source.KafkaSocketKeepaliveEnable = m.KafkaSocketKeepaliveEnable
source.consumerCloseTimeout = m.ConsumerCloseTimeout
if err := source.Open(); err != nil {
return nil, errors.Wrap(err, "opening source")
}
return source, nil
}
return m, nil
}