Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 47 additions & 26 deletions downstreamadapter/sink/kafka/helper.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
"github.com/pingcap/ticdc/pkg/sink/codec"
"github.com/pingcap/ticdc/pkg/sink/codec/common"
"github.com/pingcap/ticdc/pkg/sink/kafka"
"github.com/pingcap/ticdc/pkg/sink/kafka/claimcheck"
"github.com/pingcap/tidb/br/pkg/utils"
)

Expand All @@ -38,6 +39,7 @@
topicManager topicmanager.TopicManager
adminClient kafka.ClusterAdminClient
factory kafka.Factory
claimCheck *claimcheck.ClaimCheck
}

func (c components) close() {
Expand All @@ -47,6 +49,9 @@
if c.topicManager != nil {
c.topicManager.Close()
}
if c.claimCheck != nil {
c.claimCheck.Close()
}
}

func newKafkaSinkComponentWithFactory(ctx context.Context,
Expand All @@ -55,78 +60,94 @@
sinkConfig *config.SinkConfig,
factoryCreator kafka.FactoryCreator,
) (components, config.Protocol, error) {
kafkaComponent := components{}
var (
comp components
err error
)
// must release resources when error occurs.
defer func() {
if err != nil {
comp.close()
}
}()
protocol, err := helper.GetProtocol(utils.GetOrZero(sinkConfig.Protocol))
if err != nil {
return kafkaComponent, config.ProtocolUnknown, errors.Trace(err)
return comp, config.ProtocolUnknown, errors.Trace(err)
}

topic, err := helper.GetTopic(sinkURI)
if err != nil {
return kafkaComponent, protocol, errors.Trace(err)
return comp, protocol, errors.Trace(err)
}

options := kafka.NewOptions()
if err = options.Apply(changefeedID, sinkURI, sinkConfig); err != nil {
return kafkaComponent, protocol, errors.WrapError(errors.ErrKafkaInvalidConfig, err)
return comp, protocol, errors.WrapError(errors.ErrKafkaInvalidConfig, err)
}
options.Topic = topic

<<<<<<< HEAD

Check failure on line 89 in downstreamadapter/sink/kafka/helper.go

View workflow job for this annotation

GitHub Actions / Mac OS Build

syntax error: unexpected <<, expected }
kafkaComponent.factory, err = factoryCreator(ctx, options, changefeedID)
=======
comp.factory, err = kafka.NewSaramaFactory(ctx, options, changefeedID)
>>>>>>> bc474b549 (kafka: share one claimcheck instance across encoders (#5718))

Check failure on line 93 in downstreamadapter/sink/kafka/helper.go

View workflow job for this annotation

GitHub Actions / Mac OS Build

invalid character U+0023 '#'
if err != nil {
return kafkaComponent, protocol, errors.WrapError(errors.ErrKafkaNewProducer, err)
return comp, protocol, errors.WrapError(errors.ErrKafkaNewProducer, err)
}

<<<<<<< HEAD
kafkaComponent.eventRouter, err = eventrouter.NewEventRouter(
sinkConfig, topic, false, protocol == config.ProtocolAvro)
=======
isAvroLike := protocol == config.ProtocolAvro || protocol == config.ProtocolDebeziumAvro
comp.eventRouter, err = eventrouter.NewEventRouter(
sinkConfig, topic, false, isAvroLike)
>>>>>>> bc474b549 (kafka: share one claimcheck instance across encoders (#5718))

Check failure on line 105 in downstreamadapter/sink/kafka/helper.go

View workflow job for this annotation

GitHub Actions / Mac OS Build

invalid character U+0023 '#'
if err != nil {
return kafkaComponent, protocol, errors.Trace(err)
return comp, protocol, errors.Trace(err)
}

kafkaComponent.columnSelector, err = columnselector.New(sinkConfig)
comp.columnSelector, err = columnselector.New(sinkConfig)
if err != nil {
return kafkaComponent, protocol, errors.Trace(err)
return comp, protocol, errors.Trace(err)
}

encoderConfig, err := helper.GetEncoderConfig(changefeedID, sinkURI, protocol, sinkConfig, options.MaxMessageBytes)
if err != nil {
return kafkaComponent, protocol, errors.Trace(err)
return comp, protocol, errors.Trace(err)
}

kafkaComponent.encoderGroup, err = codec.NewEncoderGroup(ctx, sinkConfig, encoderConfig, changefeedID)
comp.claimCheck, err = claimcheck.New(ctx, encoderConfig.LargeMessageHandle, changefeedID)
if err != nil {
return kafkaComponent, protocol, errors.Trace(err)
return comp, protocol, errors.Trace(err)
}

kafkaComponent.encoder, err = codec.NewEventEncoder(ctx, encoderConfig)
comp.encoderGroup, err = codec.NewEncoderGroup(ctx, sinkConfig, encoderConfig, comp.claimCheck, changefeedID)
if err != nil {
return kafkaComponent, protocol, errors.Trace(err)
return comp, protocol, errors.Trace(err)
}

kafkaComponent.adminClient, err = kafkaComponent.factory.AdminClient(ctx)
comp.encoder, err = codec.NewEventEncoder(ctx, encoderConfig, comp.claimCheck)
if err != nil {
return kafkaComponent, protocol, errors.WrapError(errors.ErrKafkaNewProducer, err)
return comp, protocol, errors.Trace(err)
}

// We must close adminClient when this func return cause by an error
// otherwise the adminClient will never be closed and lead to a goroutine leak.
defer func() {
if err != nil && kafkaComponent.adminClient != nil {
kafkaComponent.adminClient.Close()
}
}()
comp.adminClient, err = comp.factory.AdminClient(ctx)
if err != nil {
return comp, protocol, errors.WrapError(errors.ErrKafkaNewProducer, err)
}

kafkaComponent.topicManager, err = topicmanager.GetTopicManagerAndTryCreateTopic(
comp.topicManager, err = topicmanager.GetTopicManagerAndTryCreateTopic(
ctx,
changefeedID,
topic,
options.DeriveTopicConfig(),
kafkaComponent.adminClient,
comp.adminClient,
)
if err != nil {
return kafkaComponent, protocol, errors.Trace(err)
return comp, protocol, errors.Trace(err)
}
return kafkaComponent, protocol, nil
return comp, protocol, nil
}

func newKafkaSinkComponent(
Expand Down
104 changes: 96 additions & 8 deletions downstreamadapter/sink/kafka/sink.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
"github.com/pingcap/ticdc/pkg/metrics"
"github.com/pingcap/ticdc/pkg/sink/codec/common"
"github.com/pingcap/ticdc/pkg/sink/kafka"
"github.com/pingcap/ticdc/pkg/sink/kafka/claimcheck"
"github.com/pingcap/ticdc/pkg/util"
"github.com/pingcap/ticdc/utils/chann"
"go.uber.org/atomic"
Expand Down Expand Up @@ -68,9 +69,90 @@
}

func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url.URL, sinkConfig *config.SinkConfig) error {
<<<<<<< HEAD

Check failure on line 72 in downstreamadapter/sink/kafka/sink.go

View workflow job for this annotation

GitHub Actions / Mac OS Build

syntax error: unexpected <<, expected }
comp, _, err := newKafkaSinkComponent(ctx, changefeedID, uri, sinkConfig)
defer comp.close()
return err
=======
protocol, err := helper.GetProtocol(util.GetOrZero(sinkConfig.Protocol))
if err != nil {
return errors.Trace(err)
}

topic, err := helper.GetTopic(uri)
if err != nil {
return errors.Trace(err)
}

options := kafka.NewOptions()
if err = options.Apply(changefeedID, uri, sinkConfig); err != nil {
return errors.WrapError(errors.ErrKafkaInvalidConfig, err)
}
options.Topic = topic

encoderConfig, err := helper.GetEncoderConfig(
changefeedID, uri, protocol, sinkConfig,
options.MaxMessageBytes, options.MaxBatchedBytes,
)
if err != nil {
return errors.Trace(err)
}

claimCheck, err := claimcheck.New(ctx, encoderConfig.LargeMessageHandle, changefeedID)
if err != nil {
return err
}
defer claimCheck.Close()

isAvroLike := protocol == config.ProtocolAvro || protocol == config.ProtocolDebeziumAvro
if _, err = eventrouter.NewEventRouter(sinkConfig, topic, false, isAvroLike); err != nil {
return errors.Trace(err)
}

if _, err = columnselector.New(sinkConfig); err != nil {
return errors.Trace(err)
}

factory, err := kafka.NewSaramaFactory(ctx, options, changefeedID)
if err != nil {
return errors.WrapError(errors.ErrKafkaNewProducer, err)
}

adminClient, err := factory.AdminClient(ctx)
if err != nil {
return errors.WrapError(errors.ErrKafkaNewProducer, err)
}
defer adminClient.Close()

topics, err := adminClient.GetTopicsMeta([]string{topic}, false)
if err != nil {
return errors.Trace(err)
}
if _, exists := topics[topic]; exists {
return nil
}

topicConfig := options.DeriveTopicConfig()
if !topicConfig.AutoCreate {
return errors.ErrKafkaInvalidConfig.GenWithStack("`auto-create-topic` is false, and %s not found", topic)
}

// the topic is not created, only validate.
err = adminClient.CreateTopic(&kafka.TopicDetail{
Name: topic,
NumPartitions: topicConfig.PartitionNum,
ReplicationFactor: topicConfig.ReplicationFactor,
}, true)
if err != nil {
return errors.WrapError(errors.ErrKafkaCreateTopic, err)
}

_, err = codec.NewEventEncoder(ctx, encoderConfig, claimCheck)
if err != nil {
return errors.Trace(err)
}
return nil
>>>>>>> bc474b549 (kafka: share one claimcheck instance across encoders (#5718))

Check failure on line 155 in downstreamadapter/sink/kafka/sink.go

View workflow job for this annotation

GitHub Actions / Mac OS Build

invalid character U+0023 '#'
}

func New(
Expand All @@ -89,24 +171,30 @@
protocol config.Protocol,
comp components,
) (*sink, error) {
statistics := metrics.NewStatistics(changefeedID, keyspaceID, "sink")
var (
err error
asyncProducer kafka.AsyncProducer
syncProducer kafka.SyncProducer
)
defer func() {
if err != nil {
if syncProducer != nil {
syncProducer.Close()
}
if asyncProducer != nil {
asyncProducer.Close()
}
comp.close()
if err == nil {
return
}
if syncProducer != nil {
syncProducer.Close()
}
if asyncProducer != nil {
asyncProducer.Close()
}
comp.close()
statistics.Close()
}()

<<<<<<< HEAD

Check failure on line 194 in downstreamadapter/sink/kafka/sink.go

View workflow job for this annotation

GitHub Actions / Mac OS Build

syntax error: unexpected <<, expected }
statistics := metrics.NewStatistics(changefeedID, "sink")
=======
>>>>>>> bc474b549 (kafka: share one claimcheck instance across encoders (#5718))

Check failure on line 197 in downstreamadapter/sink/kafka/sink.go

View workflow job for this annotation

GitHub Actions / Mac OS Build

invalid character U+0023 '#'
asyncProducer, err = comp.factory.AsyncProducer(ctx)
if err != nil {
return nil, err
Expand Down
75 changes: 75 additions & 0 deletions downstreamadapter/sink/kafka/sink_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,82 @@ func newKafkaSinkForTestWithProducers(ctx context.Context,
statistics := metrics.NewStatistics(changefeedID, "sink")
comp, protocol, err := newKafkaSinkComponentForTest(ctx, changefeedID, sinkURI, sinkConfig)
if err != nil {
<<<<<<< HEAD
return nil, errors.Trace(err)
=======
return nil, err
}
topic, err := helper.GetTopic(sinkURI)
if err != nil {
return nil, err
}
options := kafka.NewOptions()
if err = options.Apply(changefeedID, sinkURI, sinkConfig); err != nil {
return nil, err
}
options.Topic = topic

adminClient := kafka.NewMockClusterAdminClient(ctrl)
adminClient.EXPECT().GetTopicsMeta([]string{kafkaSinkTestTopic}, true).Return(
map[string]kafka.TopicDetail{
kafkaSinkTestTopic: {
Name: kafkaSinkTestTopic,
NumPartitions: 1,
},
}, nil)
adminClient.EXPECT().Close().AnyTimes()

metricsCollector := kafka.NewMockMetricsCollector(ctrl)
metricsCollector.EXPECT().Run(gomock.Any()).AnyTimes()

factory := kafka.NewMockFactory(ctrl)
factory.EXPECT().AsyncProducer(gomock.Any()).Return(asyncProducer, nil)
factory.EXPECT().SyncProducer(gomock.Any()).Return(syncProducer, nil)
factory.EXPECT().MetricsCollector(adminClient).Return(metricsCollector)

eventRouter, err := eventrouter.NewEventRouter(sinkConfig, topic, false, false)
if err != nil {
return nil, err
}
columnSelector, err := columnselector.New(sinkConfig)
if err != nil {
return nil, err
}
encoderConfig, err := helper.GetEncoderConfig(
changefeedID, sinkURI, protocol, sinkConfig,
options.MaxMessageBytes, options.MaxBatchedBytes,
)
if err != nil {
return nil, err
}
encoderGroup, err := codec.NewEncoderGroup(ctx, sinkConfig, encoderConfig, nil, changefeedID)
if err != nil {
return nil, err
}
encoder, err := codec.NewEventEncoder(ctx, encoderConfig, nil)
if err != nil {
return nil, err
}
topicManager, err := topicmanager.GetTopicManagerAndTryCreateTopic(
ctx,
changefeedID,
topic,
options.DeriveTopicConfig(),
adminClient,
)
if err != nil {
return nil, err
}

comp := components{
encoderGroup: encoderGroup,
encoder: encoder,
columnSelector: columnSelector,
eventRouter: eventRouter,
topicManager: topicManager,
adminClient: adminClient,
factory: factory,
>>>>>>> bc474b549 (kafka: share one claimcheck instance across encoders (#5718))
}

// We must close adminClient when this func return cause by an error
Expand Down
4 changes: 2 additions & 2 deletions downstreamadapter/sink/pulsar/helper.go
Original file line number Diff line number Diff line change
Expand Up @@ -127,12 +127,12 @@ func newPulsarSinkComponentWithFactory(ctx context.Context,
return pulsarComponent, protocol, errors.Trace(err)
}

pulsarComponent.encoderGroup, err = codec.NewEncoderGroup(ctx, sinkConfig, encoderConfig, changefeedID)
pulsarComponent.encoderGroup, err = codec.NewEncoderGroup(ctx, sinkConfig, encoderConfig, nil, changefeedID)
if err != nil {
return pulsarComponent, protocol, errors.Trace(err)
}

pulsarComponent.encoder, err = codec.NewEventEncoder(ctx, encoderConfig)
pulsarComponent.encoder, err = codec.NewEventEncoder(ctx, encoderConfig, nil)
if err != nil {
return pulsarComponent, protocol, errors.Trace(err)
}
Expand Down
2 changes: 0 additions & 2 deletions pkg/sink/codec/avro/arvo.go
Original file line number Diff line number Diff line change
Expand Up @@ -698,8 +698,6 @@ func (a *BatchEncoder) columnToAvroData(
}
}

func (a *BatchEncoder) Clean() {}

type avroEncodeResult struct {
data []byte
// header is the message header, it will be encoder into the head
Expand Down
1 change: 0 additions & 1 deletion pkg/sink/codec/bootstraper.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,6 @@ func (b *bootstrapWorker) run(ctx context.Context) error {
sendTicker := time.NewTicker(bootstrapWorkerTickerInterval)
gcTicker := time.NewTicker(bootstrapWorkerGCInterval)
defer func() {
b.rowEventEncoder.Clean()
gcTicker.Stop()
sendTicker.Stop()
}()
Expand Down
Loading
Loading