bitrix/internal/initialize/kafka.go

28 lines
608 B
Go
Raw Normal View History

2024-09-16 15:14:36 +00:00
package initialize
import (
"context"
"github.com/twmb/franz-go/pkg/kgo"
"time"
)
func KafkaConsumerInit(ctx context.Context, config Config) (*kgo.Client, error) {
kafkaClient, err := kgo.NewClient(
kgo.SeedBrokers(config.KafkaBrokers),
kgo.ConsumerGroup(config.KafkaGroup),
2025-02-25 10:13:00 +00:00
kgo.DefaultProduceTopic(config.KafkaTopicQueueBitrix),
kgo.ConsumeTopics(config.KafkaTopicQueueBitrix),
2024-09-16 15:14:36 +00:00
kgo.ConsumeResetOffset(kgo.NewOffset().AfterMilli(time.Now().UnixMilli())),
)
if err != nil {
return nil, err
}
err = kafkaClient.Ping(ctx)
if err != nil {
return nil, err
}
return kafkaClient, nil
}