forked from redpanda-data/kminion
-
Notifications
You must be signed in to change notification settings - Fork 4
/
Copy pathservice.go
101 lines (85 loc) · 2.83 KB
/
service.go
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
96
97
98
99
100
101
package kafka
import (
"context"
"fmt"
"time"
"github.com/twmb/franz-go/pkg/kerr"
"github.com/twmb/franz-go/pkg/kgo"
"github.com/twmb/franz-go/pkg/kmsg"
"github.com/twmb/franz-go/pkg/kversion"
"go.uber.org/zap"
)
type Service struct {
cfg Config
logger *zap.Logger
}
func NewService(cfg Config, logger *zap.Logger) *Service {
return &Service{
cfg: cfg,
logger: logger.Named("kafka_service"),
}
}
// CreateAndTestClient creates a client with the services default settings
// logger: will be used to log connections, errors, warnings about tls config, ...
func (s *Service) CreateAndTestClient(ctx context.Context, l *zap.Logger, opts []kgo.Opt) (*kgo.Client, error) {
logger := l.Named("kgo_client")
// Config with default options
kgoOpts, err := NewKgoConfig(s.cfg, logger)
if err != nil {
return nil, fmt.Errorf("failed to create a valid kafka Client config: %w", err)
}
// Append user (the service calling this method) provided options
kgoOpts = append(kgoOpts, opts...)
// Create kafka client
client, err := kgo.NewClient(kgoOpts...)
if err != nil {
return nil, fmt.Errorf("failed to create kafka Client: %w", err)
}
// Test connection
for {
err = s.testConnection(client, ctx)
if err == nil {
break
}
if !s.cfg.RetryInitConnection {
return nil, fmt.Errorf("failed to test connectivity to Kafka cluster %w", err)
}
logger.Warn("failed to test connectivity to Kafka cluster, retrying in 5 seconds", zap.Error(err))
time.Sleep(time.Second * 5)
}
return client, nil
}
// Brokers returns list of brokers this service is connecting to
func (s *Service) Brokers() []string {
return s.cfg.Brokers
}
// testConnection tries to fetch Broker metadata and prints some information if connection succeeds. An error will be
// returned if connecting fails.
func (s *Service) testConnection(client *kgo.Client, ctx context.Context) error {
connectCtx, cancel := context.WithTimeout(ctx, 15*time.Second)
defer cancel()
req := kmsg.MetadataRequest{
Topics: nil,
}
res, err := req.RequestWith(connectCtx, client)
if err != nil {
return fmt.Errorf("failed to request metadata: %w", err)
}
// Request versions in order to guess Kafka Cluster version
versionsReq := kmsg.NewApiVersionsRequest()
versionsRes, err := versionsReq.RequestWith(connectCtx, client)
if err != nil {
return fmt.Errorf("failed to request api versions: %w", err)
}
err = kerr.ErrorForCode(versionsRes.ErrorCode)
if err != nil {
return fmt.Errorf("failed to request api versions. Inner kafka error: %w", err)
}
versions := kversion.FromApiVersionsResponse(versionsRes)
s.logger.Debug("successfully connected to kafka cluster",
zap.Int("advertised_broker_count", len(res.Brokers)),
zap.Int("topic_count", len(res.Topics)),
zap.Int32("controller_id", res.ControllerID),
zap.String("kafka_version", versions.VersionGuess()))
return nil
}