Skip to content
Merged
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
67 changes: 39 additions & 28 deletions downstreamadapter/sink/kafka/helper.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import (
"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 @@ type components struct {
topicManager topicmanager.TopicManager
adminClient kafka.ClusterAdminClient
factory kafka.Factory
claimCheck *claimcheck.ClaimCheck
}

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

func newKafkaSinkComponent(
Expand All @@ -55,80 +60,86 @@ func newKafkaSinkComponent(
sinkURI *url.URL,
sinkConfig *config.SinkConfig,
) (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

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

isAvroLike := protocol == config.ProtocolAvro || protocol == config.ProtocolDebeziumAvro
kafkaComponent.eventRouter, err = eventrouter.NewEventRouter(
comp.eventRouter, err = eventrouter.NewEventRouter(
sinkConfig, topic, false, isAvroLike)
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, options.MaxBatchedBytes,
)
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
}
31 changes: 19 additions & 12 deletions downstreamadapter/sink/kafka/sink.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import (
"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/ticdc/pkg/util"
"github.com/pingcap/ticdc/utils/chann"
"go.uber.org/atomic"
Expand Down Expand Up @@ -95,6 +96,12 @@ func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url.
return errors.Trace(err)
}

claimCheck, err := claimcheck.New(ctx, encoderConfig.LargeMessageHandle, changefeedID)
if err != nil {
return err
}
defer claimCheck.Close()
Comment thread
3AceShowHand marked this conversation as resolved.

isAvroLike := protocol == config.ProtocolAvro || protocol == config.ProtocolDebeziumAvro
if _, err = eventrouter.NewEventRouter(sinkConfig, topic, false, isAvroLike); err != nil {
return errors.Trace(err)
Expand Down Expand Up @@ -138,12 +145,10 @@ func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url.
return errors.WrapError(errors.ErrKafkaCreateTopic, err)
}

encoder, err := codec.NewEventEncoder(ctx, encoderConfig)
_, err = codec.NewEventEncoder(ctx, encoderConfig, claimCheck)
if err != nil {
return errors.Trace(err)
}
encoder.Clean()

return nil
}

Expand All @@ -164,24 +169,26 @@ func newWithComponents(
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()
}()

statistics := metrics.NewStatistics(changefeedID, keyspaceID, "sink")
asyncProducer, err = comp.factory.AsyncProducer(ctx)
if err != nil {
return nil, err
Expand Down
4 changes: 2 additions & 2 deletions downstreamadapter/sink/kafka/sink_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -117,11 +117,11 @@ func newKafkaSinkForTestWithProducers(ctx context.Context,
if err != nil {
return nil, err
}
encoderGroup, err := codec.NewEncoderGroup(ctx, sinkConfig, encoderConfig, changefeedID)
encoderGroup, err := codec.NewEncoderGroup(ctx, sinkConfig, encoderConfig, nil, changefeedID)
if err != nil {
return nil, err
}
encoder, err := codec.NewEventEncoder(ctx, encoderConfig)
encoder, err := codec.NewEventEncoder(ctx, encoderConfig, nil)
if err != nil {
return nil, err
}
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 @@ -130,12 +130,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)
Comment thread
3AceShowHand marked this conversation as resolved.
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
9 changes: 5 additions & 4 deletions pkg/sink/codec/builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,23 +27,24 @@ import (
"github.com/pingcap/ticdc/pkg/sink/codec/debezium"
"github.com/pingcap/ticdc/pkg/sink/codec/open"
"github.com/pingcap/ticdc/pkg/sink/codec/simple"
"github.com/pingcap/ticdc/pkg/sink/kafka/claimcheck"
"go.uber.org/zap"
)

func NewEventEncoder(ctx context.Context, cfg *common.Config) (common.EventEncoder, error) {
func NewEventEncoder(ctx context.Context, cfg *common.Config, claimCheck *claimcheck.ClaimCheck) (common.EventEncoder, error) {
switch cfg.Protocol {
case config.ProtocolDefault, config.ProtocolOpen:
return open.NewBatchEncoder(ctx, cfg)
return open.NewBatchEncoder(cfg, claimCheck)
case config.ProtocolAvro:
return avro.NewAvroEncoder(ctx, cfg)
case config.ProtocolCanalJSON:
return canal.NewJSONRowEventEncoder(ctx, cfg)
return canal.NewJSONRowEventEncoder(cfg, claimCheck)
case config.ProtocolDebezium:
return debezium.NewBatchEncoder(cfg, config.GetGlobalServerConfig().ClusterID), nil
case config.ProtocolDebeziumAvro:
return debezium.NewAvroBatchEncoder(ctx, cfg, config.GetGlobalServerConfig().ClusterID)
case config.ProtocolSimple:
return simple.NewEncoder(ctx, cfg)
return simple.NewEncoder(cfg, claimCheck)
default:
return nil, errors.ErrSinkUnknownProtocol.GenWithStackByArgs(cfg.Protocol)
}
Expand Down
12 changes: 1 addition & 11 deletions pkg/sink/codec/canal/canal_json_encoder.go
Original file line number Diff line number Diff line change
Expand Up @@ -373,11 +373,7 @@ type JSONRowEventEncoder struct {
}

// NewJSONRowEventEncoder creates a new JSONRowEventEncoder
func NewJSONRowEventEncoder(ctx context.Context, config *common.Config) (common.EventEncoder, error) {
claimCheck, err := claimcheck.New(ctx, config.LargeMessageHandle, config.ChangefeedID)
if err != nil {
return nil, err
}
func NewJSONRowEventEncoder(config *common.Config, claimCheck *claimcheck.ClaimCheck) (common.EventEncoder, error) {
return &JSONRowEventEncoder{
messages: make([]*common.Message, 0, 1),
config: config,
Expand Down Expand Up @@ -582,9 +578,3 @@ func (c *JSONRowEventEncoder) EncodeDDLEvent(e *commonEvent.DDLEvent) (*common.M

return common.NewMsg(nil, value), nil
}

func (c *JSONRowEventEncoder) Clean() {
if c.claimCheck != nil {
c.claimCheck.CleanMetrics()
}
}
Loading
Loading