-
Notifications
You must be signed in to change notification settings - Fork 3
/
client.go
90 lines (68 loc) · 1.88 KB
/
client.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
package kafka
import (
"context"
"encoding/json"
"fmt"
"strings"
"github.com/rss3-network/node/internal/stream"
"github.com/rss3-network/protocol-go/schema"
"github.com/twmb/franz-go/pkg/kadm"
"github.com/twmb/franz-go/pkg/kgo"
"go.uber.org/zap"
)
type Client struct {
kafkaClient *kgo.Client
topic string
}
func New(ctx context.Context, uri, topic string) (stream.Client, error) {
brokers := strings.Split(uri, ",")
if len(brokers) == 0 {
return nil, fmt.Errorf("invalid uri: %s", uri)
}
kafkaClient, err := kgo.NewClient([]kgo.Opt{kgo.SeedBrokers(brokers...)}...)
if err != nil {
return nil, fmt.Errorf("new kafka client: %w", err)
}
kafkaAdminClient := kadm.NewClient(kafkaClient)
topicDetails, err := kafkaAdminClient.ListTopics(ctx)
if err != nil {
return nil, fmt.Errorf("list topics: %w", err)
}
if _, exists := topicDetails[topic]; !exists {
if _, err = kafkaAdminClient.CreateTopic(ctx, 1, 1, nil, topic); err != nil {
return nil, fmt.Errorf("create %s topic: %w", topic, err)
}
}
return &Client{
kafkaClient: kafkaClient,
topic: topic,
}, nil
}
func (c *Client) PushFeeds(ctx context.Context, feeds []*schema.Feed) error {
records := make([]*kgo.Record, 0, len(feeds))
for _, feed := range feeds {
record, err := c.encodeFeed(feed)
if err != nil {
return fmt.Errorf("encode feed %s: %w", feed.ID, err)
}
records = append(records, record)
}
produceResults := c.kafkaClient.ProduceSync(ctx, records...)
if err := produceResults.FirstErr(); err != nil {
return fmt.Errorf("push feeds: %w", err)
}
zap.L().Info("pushed feeds", zap.Int("records", len(records)))
return nil
}
func (c *Client) encodeFeed(feed *schema.Feed) (*kgo.Record, error) {
value, err := json.Marshal(feed)
if err != nil {
return nil, err
}
record := kgo.Record{
Topic: c.topic,
Key: []byte(feed.ID),
Value: value,
}
return &record, nil
}