From f921f6df212a939cd210a359957ecfb6b42d7132 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Mon, 24 Aug 2026 18:32:02 +0800 Subject: [PATCH 01/13] fix unknown topic and partition when get topic meta --- pkg/sink/kafka/admin.go | 3 --- pkg/sink/kafka/sarama_admin_test.go | 12 ++++-------- 2 files changed, 4 insertions(+), 11 deletions(-) diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index d45ec702d1..cc82af4f5a 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -128,9 +128,6 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool for _, meta := range metaList { if meta.Err != sarama.ErrNoError { - if meta.Err == sarama.ErrUnknownTopicOrPartition { - continue - } if !ignoreTopicError { if IsAuthorizationFailed(meta.Err) { return nil, errors.WrapError(errors.ErrKafkaAuthorizationFailed, meta.Err, "describe-topic", meta.Name) diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index 90fc3530dd..5dd42ba6bd 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -147,7 +147,7 @@ func TestGetTopicConfig(t *testing.T) { func TestGetTopicsMeta(t *testing.T) { t.Parallel() - t.Run("returns valid topics and ignores unknown topics", func(t *testing.T) { + t.Run("returns unknown topic error", func(t *testing.T) { ctrl := gomock.NewController(t) admin := NewMocksaramaClusterAdmin(ctrl) admin.EXPECT().DescribeTopics([]string{"valid-topic", "missing-topic"}).Return([]*sarama.TopicMetadata{ @@ -167,13 +167,9 @@ func TestGetTopicsMeta(t *testing.T) { topics, err := client.GetTopicsMeta([]string{"valid-topic", "missing-topic"}, false) - require.NoError(t, err) - require.Equal(t, map[string]TopicDetail{ - "valid-topic": { - Name: "valid-topic", - NumPartitions: 2, - }, - }, topics) + require.Nil(t, topics) + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, sarama.ErrUnknownTopicOrPartition) }) t.Run("missing response", func(t *testing.T) { From a4cdf2c3313b6ea43e5c9c5875a3a8c5beb15d73 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Tue, 25 Aug 2026 11:55:53 +0800 Subject: [PATCH 02/13] add unit test --- pkg/sink/kafka/sarama_admin_test.go | 33 +++++++++++++++++++++++++++-- 1 file changed, 31 insertions(+), 2 deletions(-) diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index 5dd42ba6bd..4930c51aa5 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -147,7 +147,7 @@ func TestGetTopicConfig(t *testing.T) { func TestGetTopicsMeta(t *testing.T) { t.Parallel() - t.Run("returns unknown topic error", func(t *testing.T) { + t.Run("returns unknown topic error when topic errors are not ignored", func(t *testing.T) { ctrl := gomock.NewController(t) admin := NewMocksaramaClusterAdmin(ctrl) admin.EXPECT().DescribeTopics([]string{"valid-topic", "missing-topic"}).Return([]*sarama.TopicMetadata{ @@ -172,6 +172,35 @@ func TestGetTopicsMeta(t *testing.T) { require.ErrorIs(t, err, sarama.ErrUnknownTopicOrPartition) }) + t.Run("ignores unknown topic error and returns valid topics", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"valid-topic", "missing-topic"}).Return([]*sarama.TopicMetadata{ + { + Name: "valid-topic", + Partitions: []*sarama.PartitionMetadata{{}, {}}, + }, + { + Name: "missing-topic", + Err: sarama.ErrUnknownTopicOrPartition, + }, + }, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + topics, err := client.GetTopicsMeta([]string{"valid-topic", "missing-topic"}, true) + + require.NoError(t, err) + require.Equal(t, map[string]TopicDetail{ + "valid-topic": { + Name: "valid-topic", + NumPartitions: 2, + }, + }, topics) + }) + t.Run("missing response", func(t *testing.T) { ctrl := gomock.NewController(t) admin := NewMocksaramaClusterAdmin(ctrl) @@ -244,7 +273,7 @@ func TestGetTopicsMeta(t *testing.T) { require.True(t, IsAuthorizationFailed(err)) }) - t.Run("ignored topic error", func(t *testing.T) { + t.Run("ignores non-unknown topic error", func(t *testing.T) { ctrl := gomock.NewController(t) admin := NewMocksaramaClusterAdmin(ctrl) admin.EXPECT().DescribeTopics([]string{"test-topic"}).Return([]*sarama.TopicMetadata{ From 6f8ded0312acaf8fb5582f53df7fa7b3e4a66c51 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Tue, 25 Aug 2026 15:38:13 +0800 Subject: [PATCH 03/13] Add more logs --- .../sink/topicmanager/kafka_topic_manager.go | 53 ++++++- .../topicmanager/kafka_topic_manager_test.go | 144 +++++++++++++++--- 2 files changed, 169 insertions(+), 28 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index a0110217f3..63ee631f47 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -15,6 +15,7 @@ package topicmanager import ( "context" + "fmt" "sync" "time" @@ -196,34 +197,72 @@ func (m *kafkaTopicManager) fetchAllTopicsPartitionsNum() (map[string]int32, err func (m *kafkaTopicManager) waitUntilTopicVisible( ctx context.Context, topicName string, + requiredPartitionNum int32, ) error { start := time.Now() topics := []string{topicName} + attempts := 0 + metadataFound := false + observedPartitionNum := int32(0) err := retry.Do(ctx, func() error { + attempts++ + metadataFound = false + observedPartitionNum = 0 // ignoreTopicError is set to false since we just create the topic, // make sure the topic is visible. meta, err := m.admin.GetTopicsMeta(topics, false) if err != nil { return err } - _, ok := meta[topicName] + detail, ok := meta[topicName] if !ok { return errors.ErrKafkaAdminAPI.GenWithStackByArgs("describe-topic", topicName) } + metadataFound = true + observedPartitionNum = detail.NumPartitions + if detail.NumPartitions < requiredPartitionNum { + return errors.ErrKafkaAdminAPI.GenWithStackByArgs( + "describe-topic", + fmt.Sprintf( + "%s has %d partitions, requires at least %d", + topicName, + detail.NumPartitions, + requiredPartitionNum, + ), + ) + } return nil }, retry.WithBackoffBaseDelay(500), retry.WithBackoffMaxDelay(1000), retry.WithMaxTries(6), ) if err != nil { - log.Warn("kafka topic metadata refresh failed", + if errors.Is(errors.Cause(err), context.Canceled) { + return err + } + log.Warn("kafka topic is not ready after metadata retries", zap.String("keyspace", m.changefeedID.Keyspace()), zap.String("changefeed", m.changefeedID.Name()), zap.String("topic", topicName), + zap.Int("attempts", attempts), + zap.Bool("metadataFound", metadataFound), + zap.Int32("requiredPartitionNum", requiredPartitionNum), + zap.Int32("observedPartitionNum", observedPartitionNum), zap.Duration("duration", time.Since(start)), zap.Error(err)) + return err } - return err + if attempts > 1 { + log.Info("kafka topic became ready after metadata retries", + zap.String("keyspace", m.changefeedID.Keyspace()), + zap.String("changefeed", m.changefeedID.Name()), + zap.String("topic", topicName), + zap.Int("attempts", attempts), + zap.Int32("requiredPartitionNum", requiredPartitionNum), + zap.Int32("observedPartitionNum", observedPartitionNum), + zap.Duration("duration", time.Since(start))) + } + return nil } // createTopic creates a topic with the given name @@ -260,8 +299,6 @@ func (m *kafkaTopicManager) createTopic( return 0, err } - m.tryUpdatePartitionsAndLogging(topicName, m.cfg.PartitionNum) - return m.cfg.PartitionNum, nil } @@ -304,10 +341,11 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( return 0, err } - err = m.waitUntilTopicVisible(ctx, topicName) + err = m.waitUntilTopicVisible(ctx, topicName, partitionNum) if err != nil { return 0, err } + m.tryUpdatePartitionsAndLogging(topicName, partitionNum) log.Info( "kafka topic created", @@ -331,6 +369,9 @@ func (m *kafkaTopicManager) tryStoreTopicMeta( } numPartition := detail.NumPartitions if topicName == m.defaultTopic { + if detail.NumPartitions < m.cfg.PartitionNum { + return 0, false + } numPartition = m.cfg.PartitionNum } m.tryUpdatePartitionsAndLogging(topicName, numPartition) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index e4dc2612eb..178826a1d6 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -17,6 +17,7 @@ import ( "context" "testing" + "github.com/IBM/sarama" "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/errors" @@ -26,6 +27,13 @@ import ( const kafkaTopicManagerTestTopic = "mock_topic" +func topicDetail(topic string, partitionNum int32) kafka.TopicDetail { + return kafka.TopicDetail{ + Name: topic, + NumPartitions: partitionNum, + } +} + func TestCreateTopic(t *testing.T) { t.Parallel() @@ -38,10 +46,7 @@ func TestCreateTopic(t *testing.T) { adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{kafkaTopicManagerTestTopic}, true).Return( map[string]kafka.TopicDetail{ - kafkaTopicManagerTestTopic: { - Name: kafkaTopicManagerTestTopic, - NumPartitions: 2, - }, + kafkaTopicManagerTestTopic: topicDetail(kafkaTopicManagerTestTopic, 2), }, nil) manager := newKafkaTopicManager( kafkaTopicManagerTestTopic, @@ -69,10 +74,7 @@ func TestCreateTopic(t *testing.T) { return map[string]kafka.TopicDetail{}, nil } return map[string]kafka.TopicDetail{ - createdTopic.Name: { - Name: createdTopic.Name, - NumPartitions: createdTopic.NumPartitions, - }, + createdTopic.Name: topicDetail(createdTopic.Name, createdTopic.NumPartitions), }, nil }).Times(2) adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( @@ -195,6 +197,7 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { adminClient := kafka.NewMockAdminClient(ctrl) created := false postCreateDescribeCount := 0 + var manager *kafkaTopicManager adminClient.EXPECT().GetTopicsMeta([]string{"delayed-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetTopicsMeta([]string{"delayed-topic"}, false).DoAndReturn( func([]string, bool) (map[string]kafka.TopicDetail, error) { @@ -202,16 +205,27 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { return map[string]kafka.TopicDetail{}, nil } postCreateDescribeCount++ - if postCreateDescribeCount == 1 { + _, cached := manager.topics.Load("delayed-topic") + require.False(t, cached) + switch postCreateDescribeCount { + case 1: + return nil, errors.WrapError( + errors.ErrKafkaAdminAPI, + sarama.ErrUnknownTopicOrPartition, + "describe-topic", + "delayed-topic", + ) + case 2: return map[string]kafka.TopicDetail{}, nil + case 3: + return map[string]kafka.TopicDetail{ + "delayed-topic": topicDetail("delayed-topic", 1), + }, nil } return map[string]kafka.TopicDetail{ - "delayed-topic": { - Name: "delayed-topic", - NumPartitions: 2, - }, + "delayed-topic": topicDetail("delayed-topic", 2), }, nil - }).Times(3) + }).Times(5) adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( func(detail *kafka.TopicDetail) error { require.Equal(t, &kafka.TopicDetail{ @@ -223,20 +237,109 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { return nil }) - err := EnsureTopic( - context.Background(), - common.NewChangefeedID4Test("test", "test"), + manager = newKafkaTopicManager( "delayed-topic", + common.NewChangefeedID4Test("test", "test"), + adminClient, + &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 1, + }, + ) + + partitionNum, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "delayed-topic") + + require.NoError(t, err) + require.Equal(t, int32(2), partitionNum) + require.Equal(t, 4, postCreateDescribeCount) + cachedPartitionNum, cached := manager.topics.Load("delayed-topic") + require.True(t, cached) + require.Equal(t, int32(2), cachedPartitionNum) +} + +func TestCreateTopicDoesNotCacheBeforeVisibilityRetryExhausted(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) + adminClient.EXPECT().GetTopicsMeta([]string{"never-visible-topic"}, true). + Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"never-visible-topic"}, false). + Return(map[string]kafka.TopicDetail{}, nil). + Times(7) + adminClient.EXPECT().CreateTopic(&kafka.TopicDetail{ + Name: "never-visible-topic", + NumPartitions: 2, + ReplicationFactor: 1, + }).Return(nil) + manager := newKafkaTopicManager( + "never-visible-topic", + common.NewChangefeedID4Test("test", "test"), + adminClient, &kafka.AutoCreateTopicConfig{ AutoCreate: true, PartitionNum: 2, ReplicationFactor: 1, + RequiredAcks: kafka.WaitForLocal, }, + ) + + _, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "never-visible-topic") + + require.ErrorIs(t, err, errors.ErrReachMaxTry) + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + _, cached := manager.topics.Load("never-visible-topic") + require.False(t, cached) +} + +func TestWaitUntilTopicVisibleHonorsContextCancellation(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) + ctx, cancel := context.WithCancel(context.Background()) + adminClient.EXPECT().GetTopicsMeta([]string{"cancelled-topic"}, false).DoAndReturn( + func([]string, bool) (map[string]kafka.TopicDetail, error) { + cancel() + return nil, errors.WrapError( + errors.ErrKafkaAdminAPI, + sarama.ErrUnknownTopicOrPartition, + "describe-topic", + "cancelled-topic", + ) + }) + manager := newKafkaTopicManager( + "cancelled-topic", + common.NewChangefeedID4Test("test", "test"), + adminClient, + &kafka.AutoCreateTopicConfig{PartitionNum: 2}, + ) + + err := manager.waitUntilTopicVisible(ctx, "cancelled-topic", 2) + + require.ErrorIs(t, err, context.Canceled) +} + +func TestWaitUntilTopicVisibleAllowsAdditionalPartitions(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) + adminClient.EXPECT().GetTopicsMeta([]string{"expanded-topic"}, false).Return( + map[string]kafka.TopicDetail{ + "expanded-topic": topicDetail("expanded-topic", 3), + }, nil) + manager := newKafkaTopicManager( + "expanded-topic", + common.NewChangefeedID4Test("test", "test"), adminClient, + &kafka.AutoCreateTopicConfig{PartitionNum: 2}, ) + err := manager.waitUntilTopicVisible(context.Background(), "expanded-topic", 2) + require.NoError(t, err) - require.Equal(t, 2, postCreateDescribeCount) } func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { @@ -246,10 +349,7 @@ func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{"existing-topic"}, true).Return( map[string]kafka.TopicDetail{ - "existing-topic": { - Name: "existing-topic", - NumPartitions: 2, - }, + "existing-topic": topicDetail("existing-topic", 2), }, nil) manager, err := GetTopicManagerAndTryCreateTopic( From c4eb2927d0d439ad5c28d66e5f73abc558401579 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Tue, 25 Aug 2026 15:52:56 +0800 Subject: [PATCH 04/13] add more log --- .../sink/topicmanager/kafka_topic_manager.go | 20 +++++----- .../topicmanager/kafka_topic_manager_test.go | 27 ++++++++++++++ pkg/sink/kafka/admin.go | 19 ++++++++++ pkg/sink/kafka/sarama_admin_test.go | 37 +++++++++++++++++++ 4 files changed, 92 insertions(+), 11 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 63ee631f47..3e1646b7a1 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -15,7 +15,6 @@ package topicmanager import ( "context" - "fmt" "sync" "time" @@ -221,26 +220,25 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( metadataFound = true observedPartitionNum = detail.NumPartitions if detail.NumPartitions < requiredPartitionNum { - return errors.ErrKafkaAdminAPI.GenWithStackByArgs( - "describe-topic", - fmt.Sprintf( - "%s has %d partitions, requires at least %d", - topicName, - detail.NumPartitions, - requiredPartitionNum, - ), - ) + return errors.ErrKafkaAdminAPI.GenWithStackByArgs("describe-topic", topicName) } return nil }, retry.WithBackoffBaseDelay(500), retry.WithBackoffMaxDelay(1000), retry.WithMaxTries(6), + retry.WithIsRetryableErr(func(err error) bool { + // A direct ErrKafkaAdminAPI is generated above when the topic metadata + // is missing or has too few partitions. Admin client errors wrap their + // original cause and are classified by Kafka error semantics. + return errors.ErrKafkaAdminAPI.Equal(errors.Cause(err)) || + kafka.IsRetryableTopicMetadataError(err) + }), ) if err != nil { if errors.Is(errors.Cause(err), context.Canceled) { return err } - log.Warn("kafka topic is not ready after metadata retries", + log.Warn("kafka topic readiness check failed", zap.String("keyspace", m.changefeedID.Keyspace()), zap.String("changefeed", m.changefeedID.Name()), zap.String("topic", topicName), diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 178826a1d6..bc8b20063c 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -321,6 +321,33 @@ func TestWaitUntilTopicVisibleHonorsContextCancellation(t *testing.T) { require.ErrorIs(t, err, context.Canceled) } +func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) + adminClient.EXPECT().GetTopicsMeta([]string{"invalid-topic"}, false).Return( + nil, + errors.WrapError( + errors.ErrKafkaAdminAPI, + sarama.ErrInvalidTopic, + "describe-topic", + "invalid-topic", + ), + ).Times(1) + manager := newKafkaTopicManager( + "invalid-topic", + common.NewChangefeedID4Test("test", "test"), + adminClient, + &kafka.AutoCreateTopicConfig{PartitionNum: 2}, + ) + + err := manager.waitUntilTopicVisible(context.Background(), "invalid-topic", 2) + + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, sarama.ErrInvalidTopic) +} + func TestWaitUntilTopicVisibleAllowsAdditionalPartitions(t *testing.T) { t.Parallel() diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index cc82af4f5a..93f79c0724 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -156,6 +156,25 @@ func IsAuthorizationFailed(err error) bool { errors.Is(err, sarama.ErrClusterAuthorizationFailed) } +// IsRetryableTopicMetadataError reports whether a Kafka metadata error can be +// caused by a temporary topic, broker, network, or controller state. +func IsRetryableTopicMetadataError(err error) bool { + return errors.Is(err, sarama.ErrUnknownTopicOrPartition) || + errors.Is(err, sarama.ErrLeaderNotAvailable) || + errors.Is(err, sarama.ErrNotLeaderForPartition) || + errors.Is(err, sarama.ErrRequestTimedOut) || + errors.Is(err, sarama.ErrBrokerNotAvailable) || + errors.Is(err, sarama.ErrReplicaNotAvailable) || + errors.Is(err, sarama.ErrStaleControllerEpochCode) || + errors.Is(err, sarama.ErrNetworkException) || + errors.Is(err, sarama.ErrNotController) || + errors.Is(err, sarama.ErrKafkaStorageError) || + errors.Is(err, sarama.ErrOutOfBrokers) || + errors.Is(err, sarama.ErrBrokerNotFound) || + errors.Is(err, sarama.ErrIncompleteResponse) || + errors.Is(err, sarama.ErrControllerNotAvailable) +} + func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string]int32, error) { result := make(map[string]int32, len(topics)) for _, topic := range topics { diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index 4930c51aa5..6451ca8d0d 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -312,6 +312,43 @@ func TestIsAuthorizationFailed(t *testing.T) { } } +func TestIsRetryableTopicMetadataError(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + err error + expected bool + }{ + {name: "unknown topic", err: sarama.ErrUnknownTopicOrPartition, expected: true}, + {name: "leader unavailable", err: sarama.ErrLeaderNotAvailable, expected: true}, + {name: "request timeout", err: sarama.ErrRequestTimedOut, expected: true}, + {name: "network exception", err: sarama.ErrNetworkException, expected: true}, + {name: "controller changed", err: sarama.ErrNotController, expected: true}, + {name: "no broker available", err: sarama.ErrOutOfBrokers, expected: true}, + { + name: "wrapped retryable error", + err: errors.WrapError( + errors.ErrKafkaAdminAPI, + sarama.ErrUnknownTopicOrPartition, + "describe-topic", + "test-topic", + ), + expected: true, + }, + {name: "authorization failure", err: sarama.ErrTopicAuthorizationFailed}, + {name: "invalid topic", err: sarama.ErrInvalidTopic}, + {name: "unknown broker error", err: sarama.ErrUnknown}, + {name: "context cancellation", err: context.Canceled}, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + require.Equal(t, test.expected, IsRetryableTopicMetadataError(test.err)) + }) + } +} + func TestCreateTopic(t *testing.T) { t.Parallel() From 4d99d668e0ca77923a88d0db55ea260e8a743c60 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Tue, 25 Aug 2026 16:03:06 +0800 Subject: [PATCH 05/13] no need to compare partition num --- .../sink/topicmanager/kafka_topic_manager.go | 15 ++------- .../topicmanager/kafka_topic_manager_test.go | 33 +++---------------- 2 files changed, 7 insertions(+), 41 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 3e1646b7a1..4a283e640e 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -196,7 +196,6 @@ func (m *kafkaTopicManager) fetchAllTopicsPartitionsNum() (map[string]int32, err func (m *kafkaTopicManager) waitUntilTopicVisible( ctx context.Context, topicName string, - requiredPartitionNum int32, ) error { start := time.Now() topics := []string{topicName} @@ -219,17 +218,14 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( } metadataFound = true observedPartitionNum = detail.NumPartitions - if detail.NumPartitions < requiredPartitionNum { - return errors.ErrKafkaAdminAPI.GenWithStackByArgs("describe-topic", topicName) - } return nil }, retry.WithBackoffBaseDelay(500), retry.WithBackoffMaxDelay(1000), retry.WithMaxTries(6), retry.WithIsRetryableErr(func(err error) bool { // A direct ErrKafkaAdminAPI is generated above when the topic metadata - // is missing or has too few partitions. Admin client errors wrap their - // original cause and are classified by Kafka error semantics. + // is missing. Admin client errors wrap their original cause and are + // classified by Kafka error semantics. return errors.ErrKafkaAdminAPI.Equal(errors.Cause(err)) || kafka.IsRetryableTopicMetadataError(err) }), @@ -244,7 +240,6 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( zap.String("topic", topicName), zap.Int("attempts", attempts), zap.Bool("metadataFound", metadataFound), - zap.Int32("requiredPartitionNum", requiredPartitionNum), zap.Int32("observedPartitionNum", observedPartitionNum), zap.Duration("duration", time.Since(start)), zap.Error(err)) @@ -256,7 +251,6 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( zap.String("changefeed", m.changefeedID.Name()), zap.String("topic", topicName), zap.Int("attempts", attempts), - zap.Int32("requiredPartitionNum", requiredPartitionNum), zap.Int32("observedPartitionNum", observedPartitionNum), zap.Duration("duration", time.Since(start))) } @@ -339,7 +333,7 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( return 0, err } - err = m.waitUntilTopicVisible(ctx, topicName, partitionNum) + err = m.waitUntilTopicVisible(ctx, topicName) if err != nil { return 0, err } @@ -367,9 +361,6 @@ func (m *kafkaTopicManager) tryStoreTopicMeta( } numPartition := detail.NumPartitions if topicName == m.defaultTopic { - if detail.NumPartitions < m.cfg.PartitionNum { - return 0, false - } numPartition = m.cfg.PartitionNum } m.tryUpdatePartitionsAndLogging(topicName, numPartition) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index bc8b20063c..5a624bd1ec 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -217,15 +217,11 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { ) case 2: return map[string]kafka.TopicDetail{}, nil - case 3: - return map[string]kafka.TopicDetail{ - "delayed-topic": topicDetail("delayed-topic", 1), - }, nil } return map[string]kafka.TopicDetail{ "delayed-topic": topicDetail("delayed-topic", 2), }, nil - }).Times(5) + }).Times(4) adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( func(detail *kafka.TopicDetail) error { require.Equal(t, &kafka.TopicDetail{ @@ -252,7 +248,7 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { require.NoError(t, err) require.Equal(t, int32(2), partitionNum) - require.Equal(t, 4, postCreateDescribeCount) + require.Equal(t, 3, postCreateDescribeCount) cachedPartitionNum, cached := manager.topics.Load("delayed-topic") require.True(t, cached) require.Equal(t, int32(2), cachedPartitionNum) @@ -316,7 +312,7 @@ func TestWaitUntilTopicVisibleHonorsContextCancellation(t *testing.T) { &kafka.AutoCreateTopicConfig{PartitionNum: 2}, ) - err := manager.waitUntilTopicVisible(ctx, "cancelled-topic", 2) + err := manager.waitUntilTopicVisible(ctx, "cancelled-topic") require.ErrorIs(t, err, context.Canceled) } @@ -342,33 +338,12 @@ func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { &kafka.AutoCreateTopicConfig{PartitionNum: 2}, ) - err := manager.waitUntilTopicVisible(context.Background(), "invalid-topic", 2) + err := manager.waitUntilTopicVisible(context.Background(), "invalid-topic") require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) require.ErrorIs(t, err, sarama.ErrInvalidTopic) } -func TestWaitUntilTopicVisibleAllowsAdditionalPartitions(t *testing.T) { - t.Parallel() - - ctrl := gomock.NewController(t) - adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"expanded-topic"}, false).Return( - map[string]kafka.TopicDetail{ - "expanded-topic": topicDetail("expanded-topic", 3), - }, nil) - manager := newKafkaTopicManager( - "expanded-topic", - common.NewChangefeedID4Test("test", "test"), - adminClient, - &kafka.AutoCreateTopicConfig{PartitionNum: 2}, - ) - - err := manager.waitUntilTopicVisible(context.Background(), "expanded-topic", 2) - - require.NoError(t, err) -} - func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { t.Parallel() From 4d3bab167a2f5fd6eaa322f3b98b05abaa25929f Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Tue, 25 Aug 2026 16:49:28 +0800 Subject: [PATCH 06/13] remove useless logs --- .../sink/topicmanager/kafka_topic_manager.go | 30 ++----------------- 1 file changed, 3 insertions(+), 27 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 4a283e640e..673528c380 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -199,25 +199,17 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( ) error { start := time.Now() topics := []string{topicName} - attempts := 0 - metadataFound := false - observedPartitionNum := int32(0) err := retry.Do(ctx, func() error { - attempts++ - metadataFound = false - observedPartitionNum = 0 // ignoreTopicError is set to false since we just create the topic, // make sure the topic is visible. meta, err := m.admin.GetTopicsMeta(topics, false) if err != nil { return err } - detail, ok := meta[topicName] + _, ok := meta[topicName] if !ok { return errors.ErrKafkaAdminAPI.GenWithStackByArgs("describe-topic", topicName) } - metadataFound = true - observedPartitionNum = detail.NumPartitions return nil }, retry.WithBackoffBaseDelay(500), retry.WithBackoffMaxDelay(1000), @@ -231,30 +223,14 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( }), ) if err != nil { - if errors.Is(errors.Cause(err), context.Canceled) { - return err - } - log.Warn("kafka topic readiness check failed", + log.Warn("kafka topic metadata refresh failed", zap.String("keyspace", m.changefeedID.Keyspace()), zap.String("changefeed", m.changefeedID.Name()), zap.String("topic", topicName), - zap.Int("attempts", attempts), - zap.Bool("metadataFound", metadataFound), - zap.Int32("observedPartitionNum", observedPartitionNum), zap.Duration("duration", time.Since(start)), zap.Error(err)) - return err } - if attempts > 1 { - log.Info("kafka topic became ready after metadata retries", - zap.String("keyspace", m.changefeedID.Keyspace()), - zap.String("changefeed", m.changefeedID.Name()), - zap.String("topic", topicName), - zap.Int("attempts", attempts), - zap.Int32("observedPartitionNum", observedPartitionNum), - zap.Duration("duration", time.Since(start))) - } - return nil + return err } // createTopic creates a topic with the given name From 4d8e9b1f1ed70342a8863a4ea88b16b0bbbdc7e2 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Tue, 25 Aug 2026 16:54:25 +0800 Subject: [PATCH 07/13] decouple test from sarama --- .../topicmanager/kafka_topic_manager_test.go | 32 ++++--------------- 1 file changed, 6 insertions(+), 26 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 5a624bd1ec..868ba07926 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -17,7 +17,6 @@ import ( "context" "testing" - "github.com/IBM/sarama" "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/errors" @@ -207,21 +206,13 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { postCreateDescribeCount++ _, cached := manager.topics.Load("delayed-topic") require.False(t, cached) - switch postCreateDescribeCount { - case 1: - return nil, errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrUnknownTopicOrPartition, - "describe-topic", - "delayed-topic", - ) - case 2: + if postCreateDescribeCount == 1 { return map[string]kafka.TopicDetail{}, nil } return map[string]kafka.TopicDetail{ "delayed-topic": topicDetail("delayed-topic", 2), }, nil - }).Times(4) + }).Times(3) adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( func(detail *kafka.TopicDetail) error { require.Equal(t, &kafka.TopicDetail{ @@ -248,7 +239,7 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { require.NoError(t, err) require.Equal(t, int32(2), partitionNum) - require.Equal(t, 3, postCreateDescribeCount) + require.Equal(t, 2, postCreateDescribeCount) cachedPartitionNum, cached := manager.topics.Load("delayed-topic") require.True(t, cached) require.Equal(t, int32(2), cachedPartitionNum) @@ -298,12 +289,7 @@ func TestWaitUntilTopicVisibleHonorsContextCancellation(t *testing.T) { adminClient.EXPECT().GetTopicsMeta([]string{"cancelled-topic"}, false).DoAndReturn( func([]string, bool) (map[string]kafka.TopicDetail, error) { cancel() - return nil, errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrUnknownTopicOrPartition, - "describe-topic", - "cancelled-topic", - ) + return map[string]kafka.TopicDetail{}, nil }) manager := newKafkaTopicManager( "cancelled-topic", @@ -324,12 +310,7 @@ func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{"invalid-topic"}, false).Return( nil, - errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrInvalidTopic, - "describe-topic", - "invalid-topic", - ), + errors.ErrKafkaInvalidConfig.GenWithStack("invalid topic"), ).Times(1) manager := newKafkaTopicManager( "invalid-topic", @@ -340,8 +321,7 @@ func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { err := manager.waitUntilTopicVisible(context.Background(), "invalid-topic") - require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) - require.ErrorIs(t, err, sarama.ErrInvalidTopic) + require.ErrorIs(t, err, errors.ErrKafkaInvalidConfig) } func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { From f6a82005fdee289726b14c850845b8878685b766 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Tue, 25 Aug 2026 20:37:08 +0800 Subject: [PATCH 08/13] use unretryable error --- .../sink/topicmanager/kafka_topic_manager.go | 6 +- .../topicmanager/kafka_topic_manager_test.go | 151 ++++-------------- pkg/sink/kafka/admin.go | 34 ++-- pkg/sink/kafka/sarama_admin_test.go | 51 ++++-- 4 files changed, 85 insertions(+), 157 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 673528c380..5dc68119e4 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -215,11 +215,7 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( retry.WithBackoffMaxDelay(1000), retry.WithMaxTries(6), retry.WithIsRetryableErr(func(err error) bool { - // A direct ErrKafkaAdminAPI is generated above when the topic metadata - // is missing. Admin client errors wrap their original cause and are - // classified by Kafka error semantics. - return errors.ErrKafkaAdminAPI.Equal(errors.Cause(err)) || - kafka.IsRetryableTopicMetadataError(err) + return !kafka.IsUnretryableTopicMetadataError(err) }), ) if err != nil { diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 868ba07926..0976339f61 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -15,8 +15,10 @@ package topicmanager import ( "context" + "io" "testing" + "github.com/IBM/sarama" "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/errors" @@ -66,23 +68,41 @@ func TestCreateTopic(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) var createdTopic *kafka.TopicDetail + postCreateDescribeCount := 0 + var manager *kafkaTopicManager adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).DoAndReturn( func([]string, bool) (map[string]kafka.TopicDetail, error) { if createdTopic == nil { - return map[string]kafka.TopicDetail{}, nil + return nil, errors.WrapError( + errors.ErrKafkaAdminAPI, + sarama.ErrUnknownTopicOrPartition, + "describe-topic", + "new-topic", + ) + } + postCreateDescribeCount++ + _, cached := manager.topics.Load("new-topic") + require.False(t, cached) + if postCreateDescribeCount == 1 { + return nil, errors.WrapError( + errors.ErrKafkaAdminAPI, + io.EOF, + "describe-topic", + "new-topic", + ) } return map[string]kafka.TopicDetail{ createdTopic.Name: topicDetail(createdTopic.Name, createdTopic.NumPartitions), }, nil - }).Times(2) + }).Times(3) adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( func(detail *kafka.TopicDetail) error { copy := *detail createdTopic = © return nil }) - manager := newKafkaTopicManager( + manager = newKafkaTopicManager( kafkaTopicManagerTestTopic, changefeedID, adminClient, @@ -103,6 +123,7 @@ func TestCreateTopic(t *testing.T) { NumPartitions: 2, ReplicationFactor: 1, }, createdTopic) + require.Equal(t, 2, postCreateDescribeCount) partitionsNum, err := manager.GetPartitionNum(context.Background(), "new-topic") require.NoError(t, err) require.Equal(t, int32(2), partitionsNum) @@ -189,120 +210,6 @@ func TestCreateTopicValidatesReplicationFactor(t *testing.T) { require.ErrorContains(t, err, "`replication-factor` 1 is smaller than the `min.insync.replicas` 2 of broker") } -func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { - t.Parallel() - - ctrl := gomock.NewController(t) - adminClient := kafka.NewMockAdminClient(ctrl) - created := false - postCreateDescribeCount := 0 - var manager *kafkaTopicManager - adminClient.EXPECT().GetTopicsMeta([]string{"delayed-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) - adminClient.EXPECT().GetTopicsMeta([]string{"delayed-topic"}, false).DoAndReturn( - func([]string, bool) (map[string]kafka.TopicDetail, error) { - if !created { - return map[string]kafka.TopicDetail{}, nil - } - postCreateDescribeCount++ - _, cached := manager.topics.Load("delayed-topic") - require.False(t, cached) - if postCreateDescribeCount == 1 { - return map[string]kafka.TopicDetail{}, nil - } - return map[string]kafka.TopicDetail{ - "delayed-topic": topicDetail("delayed-topic", 2), - }, nil - }).Times(3) - adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( - func(detail *kafka.TopicDetail) error { - require.Equal(t, &kafka.TopicDetail{ - Name: "delayed-topic", - NumPartitions: 2, - ReplicationFactor: 1, - }, detail) - created = true - return nil - }) - - manager = newKafkaTopicManager( - "delayed-topic", - common.NewChangefeedID4Test("test", "test"), - adminClient, - &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, - ReplicationFactor: 1, - }, - ) - - partitionNum, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "delayed-topic") - - require.NoError(t, err) - require.Equal(t, int32(2), partitionNum) - require.Equal(t, 2, postCreateDescribeCount) - cachedPartitionNum, cached := manager.topics.Load("delayed-topic") - require.True(t, cached) - require.Equal(t, int32(2), cachedPartitionNum) -} - -func TestCreateTopicDoesNotCacheBeforeVisibilityRetryExhausted(t *testing.T) { - t.Parallel() - - ctrl := gomock.NewController(t) - adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"never-visible-topic"}, true). - Return(map[string]kafka.TopicDetail{}, nil) - adminClient.EXPECT().GetTopicsMeta([]string{"never-visible-topic"}, false). - Return(map[string]kafka.TopicDetail{}, nil). - Times(7) - adminClient.EXPECT().CreateTopic(&kafka.TopicDetail{ - Name: "never-visible-topic", - NumPartitions: 2, - ReplicationFactor: 1, - }).Return(nil) - manager := newKafkaTopicManager( - "never-visible-topic", - common.NewChangefeedID4Test("test", "test"), - adminClient, - &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, - ReplicationFactor: 1, - RequiredAcks: kafka.WaitForLocal, - }, - ) - - _, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "never-visible-topic") - - require.ErrorIs(t, err, errors.ErrReachMaxTry) - require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) - _, cached := manager.topics.Load("never-visible-topic") - require.False(t, cached) -} - -func TestWaitUntilTopicVisibleHonorsContextCancellation(t *testing.T) { - t.Parallel() - - ctrl := gomock.NewController(t) - adminClient := kafka.NewMockAdminClient(ctrl) - ctx, cancel := context.WithCancel(context.Background()) - adminClient.EXPECT().GetTopicsMeta([]string{"cancelled-topic"}, false).DoAndReturn( - func([]string, bool) (map[string]kafka.TopicDetail, error) { - cancel() - return map[string]kafka.TopicDetail{}, nil - }) - manager := newKafkaTopicManager( - "cancelled-topic", - common.NewChangefeedID4Test("test", "test"), - adminClient, - &kafka.AutoCreateTopicConfig{PartitionNum: 2}, - ) - - err := manager.waitUntilTopicVisible(ctx, "cancelled-topic") - - require.ErrorIs(t, err, context.Canceled) -} - func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { t.Parallel() @@ -310,7 +217,12 @@ func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{"invalid-topic"}, false).Return( nil, - errors.ErrKafkaInvalidConfig.GenWithStack("invalid topic"), + errors.WrapError( + errors.ErrKafkaAdminAPI, + sarama.ErrInvalidTopic, + "describe-topic", + "invalid-topic", + ), ).Times(1) manager := newKafkaTopicManager( "invalid-topic", @@ -321,7 +233,8 @@ func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { err := manager.waitUntilTopicVisible(context.Background(), "invalid-topic") - require.ErrorIs(t, err, errors.ErrKafkaInvalidConfig) + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, sarama.ErrInvalidTopic) } func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index 93f79c0724..dfb61ff05f 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -156,23 +156,23 @@ func IsAuthorizationFailed(err error) bool { errors.Is(err, sarama.ErrClusterAuthorizationFailed) } -// IsRetryableTopicMetadataError reports whether a Kafka metadata error can be -// caused by a temporary topic, broker, network, or controller state. -func IsRetryableTopicMetadataError(err error) bool { - return errors.Is(err, sarama.ErrUnknownTopicOrPartition) || - errors.Is(err, sarama.ErrLeaderNotAvailable) || - errors.Is(err, sarama.ErrNotLeaderForPartition) || - errors.Is(err, sarama.ErrRequestTimedOut) || - errors.Is(err, sarama.ErrBrokerNotAvailable) || - errors.Is(err, sarama.ErrReplicaNotAvailable) || - errors.Is(err, sarama.ErrStaleControllerEpochCode) || - errors.Is(err, sarama.ErrNetworkException) || - errors.Is(err, sarama.ErrNotController) || - errors.Is(err, sarama.ErrKafkaStorageError) || - errors.Is(err, sarama.ErrOutOfBrokers) || - errors.Is(err, sarama.ErrBrokerNotFound) || - errors.Is(err, sarama.ErrIncompleteResponse) || - errors.Is(err, sarama.ErrControllerNotAvailable) +// IsUnretryableTopicMetadataError reports whether a Kafka metadata request +// requires a configuration, credential, permission, or request change to succeed. +func IsUnretryableTopicMetadataError(err error) bool { + if IsAuthorizationFailed(err) || + errors.Is(err, errors.ErrKafkaInvalidConfig) || + errors.Is(err, sarama.ErrInvalidTopic) || + errors.Is(err, sarama.ErrInvalidConfig) || + errors.Is(err, sarama.ErrSASLAuthenticationFailed) || + errors.Is(err, sarama.ErrUnsupportedSASLMechanism) || + errors.Is(err, sarama.ErrIllegalSASLState) || + errors.Is(err, sarama.ErrUnsupportedVersion) || + errors.Is(err, sarama.ErrInvalidRequest) { + return true + } + + var configErr sarama.ConfigurationError + return errors.As(err, &configErr) } func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string]int32, error) { diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index 6451ca8d0d..efb34ea361 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -312,39 +312,58 @@ func TestIsAuthorizationFailed(t *testing.T) { } } -func TestIsRetryableTopicMetadataError(t *testing.T) { +func TestIsUnretryableTopicMetadataError(t *testing.T) { t.Parallel() tests := []struct { - name string - err error - expected bool + name string + err error + unretryable bool }{ - {name: "unknown topic", err: sarama.ErrUnknownTopicOrPartition, expected: true}, - {name: "leader unavailable", err: sarama.ErrLeaderNotAvailable, expected: true}, - {name: "request timeout", err: sarama.ErrRequestTimedOut, expected: true}, - {name: "network exception", err: sarama.ErrNetworkException, expected: true}, - {name: "controller changed", err: sarama.ErrNotController, expected: true}, - {name: "no broker available", err: sarama.ErrOutOfBrokers, expected: true}, + {name: "unknown topic", err: sarama.ErrUnknownTopicOrPartition}, + {name: "leader unavailable", err: sarama.ErrLeaderNotAvailable}, + {name: "request timeout", err: sarama.ErrRequestTimedOut}, + {name: "network exception", err: sarama.ErrNetworkException}, + {name: "controller changed", err: sarama.ErrNotController}, + {name: "no broker available", err: sarama.ErrOutOfBrokers}, + {name: "EOF", err: io.EOF}, + {name: "unknown broker error", err: sarama.ErrUnknown}, { - name: "wrapped retryable error", + name: "wrapped unknown topic", err: errors.WrapError( errors.ErrKafkaAdminAPI, sarama.ErrUnknownTopicOrPartition, "describe-topic", "test-topic", ), - expected: true, }, - {name: "authorization failure", err: sarama.ErrTopicAuthorizationFailed}, - {name: "invalid topic", err: sarama.ErrInvalidTopic}, - {name: "unknown broker error", err: sarama.ErrUnknown}, {name: "context cancellation", err: context.Canceled}, + {name: "TiCDC invalid config", err: errors.ErrKafkaInvalidConfig.GenWithStack("invalid config"), unretryable: true}, + {name: "topic authorization failure", err: sarama.ErrTopicAuthorizationFailed, unretryable: true}, + {name: "cluster authorization failure", err: sarama.ErrClusterAuthorizationFailed, unretryable: true}, + {name: "invalid topic", err: sarama.ErrInvalidTopic, unretryable: true}, + {name: "invalid config", err: sarama.ErrInvalidConfig, unretryable: true}, + {name: "SASL authentication failure", err: sarama.ErrSASLAuthenticationFailed, unretryable: true}, + {name: "unsupported SASL mechanism", err: sarama.ErrUnsupportedSASLMechanism, unretryable: true}, + {name: "illegal SASL state", err: sarama.ErrIllegalSASLState, unretryable: true}, + {name: "unsupported version", err: sarama.ErrUnsupportedVersion, unretryable: true}, + {name: "invalid request", err: sarama.ErrInvalidRequest, unretryable: true}, + {name: "client configuration error", err: sarama.ConfigurationError("invalid client config"), unretryable: true}, + { + name: "wrapped invalid topic", + err: errors.WrapError( + errors.ErrKafkaAdminAPI, + sarama.ErrInvalidTopic, + "describe-topic", + "test-topic", + ), + unretryable: true, + }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { - require.Equal(t, test.expected, IsRetryableTopicMetadataError(test.err)) + require.Equal(t, test.unretryable, IsUnretryableTopicMetadataError(test.err)) }) } } From b578159cfb0967a9d2d02c8efe0db3ea2a7b926d Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Wed, 26 Aug 2026 12:05:30 +0800 Subject: [PATCH 09/13] do not push if the sink is closed --- downstreamadapter/sink/kafka/sink.go | 5 +++++ downstreamadapter/sink/kafka/sink_test.go | 15 +++++++++++++++ 2 files changed, 20 insertions(+) diff --git a/downstreamadapter/sink/kafka/sink.go b/downstreamadapter/sink/kafka/sink.go index caf9f324cb..5faf5f3296 100644 --- a/downstreamadapter/sink/kafka/sink.go +++ b/downstreamadapter/sink/kafka/sink.go @@ -232,6 +232,9 @@ func (s *sink) IsNormal() bool { } func (s *sink) AddDMLEvent(event *commonEvent.DMLEvent) { + if !s.isNormal.Load() { + return + } s.eventChan.Push(event) } @@ -567,6 +570,8 @@ func (s *sink) getAllTableNames(ts uint64) []*commonEvent.SchemaTableName { } func (s *sink) Close() { + s.isNormal.Store(false) + s.close() s.ddlProducer.Close() s.dmlProducer.Close() s.comp.close() diff --git a/downstreamadapter/sink/kafka/sink_test.go b/downstreamadapter/sink/kafka/sink_test.go index 280bdf3075..e518a35609 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -313,8 +313,23 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { require.NoError(t, err) require.Zero(t, closeCount.Load()) + require.True(t, kafkaSink.IsNormal()) + + kafkaSink.isNormal.Store(false) + kafkaSink.AddDMLEvent(&commonEvent.DMLEvent{}) + require.Zero(t, kafkaSink.eventChan.Len()) + kafkaSink.isNormal.Store(true) + kafkaSink.Close() require.Equal(t, int64(4), closeCount.Load()) + require.False(t, kafkaSink.IsNormal()) + + _, ok, err := kafkaSink.eventChan.GetWithContext(t.Context()) + require.NoError(t, err) + require.False(t, ok) + _, ok, err = kafkaSink.rowChan.GetWithContext(t.Context()) + require.NoError(t, err) + require.False(t, ok) }) } From 4a427036f9a2217d576dbc2c8ec65ecf556cdfdc Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Wed, 26 Aug 2026 14:33:27 +0800 Subject: [PATCH 10/13] share the client between all sarama components --- downstreamadapter/sink/kafka/helper.go | 6 +- downstreamadapter/sink/kafka/sink.go | 15 ++++- downstreamadapter/sink/kafka/sink_test.go | 31 +++++----- pkg/sink/kafka/admin.go | 37 +++--------- pkg/sink/kafka/sarama_admin_mock.go | 66 --------------------- pkg/sink/kafka/sarama_admin_test.go | 44 +++----------- pkg/sink/kafka/sarama_async_producer.go | 26 +------- pkg/sink/kafka/sarama_factory.go | 47 +++++---------- pkg/sink/kafka/sarama_sync_producer.go | 17 ------ pkg/sink/kafka/sarama_sync_producer_mock.go | 51 ---------------- pkg/sink/kafka/sarama_sync_producer_test.go | 39 +++--------- 11 files changed, 72 insertions(+), 307 deletions(-) diff --git a/downstreamadapter/sink/kafka/helper.go b/downstreamadapter/sink/kafka/helper.go index fbbdd4cd50..b3353afe25 100644 --- a/downstreamadapter/sink/kafka/helper.go +++ b/downstreamadapter/sink/kafka/helper.go @@ -42,12 +42,12 @@ type components struct { } func (c components) close() { - if c.adminClient != nil { - c.adminClient.Close() - } if c.topicManager != nil { c.topicManager.Close() } + if c.adminClient != nil { + c.adminClient.Close() + } if c.claimCheck != nil { c.claimCheck.Close() } diff --git a/downstreamadapter/sink/kafka/sink.go b/downstreamadapter/sink/kafka/sink.go index 5faf5f3296..a156f38021 100644 --- a/downstreamadapter/sink/kafka/sink.go +++ b/downstreamadapter/sink/kafka/sink.go @@ -168,13 +168,15 @@ func newWithComponents( if err == nil { return } + // Closing the admin releases the shared Sarama client and unblocks + // producer shutdown when Kafka is unhealthy. + comp.close() if syncProducer != nil { syncProducer.Close() } if asyncProducer != nil { asyncProducer.Close() } - comp.close() statistics.Close() }() @@ -263,6 +265,7 @@ func (s *sink) WriteBlockEvent(event commonEvent.BlockEvent) error { } func (s *sink) close() { + s.isNormal.Store(false) s.eventChan.Close() s.rowChan.Close() } @@ -318,6 +321,11 @@ func (s *sink) calculateKeyPartitions(ctx context.Context) error { if err != nil { return err } + select { + case <-ctx.Done(): + return context.Cause(ctx) + default: + } s.rowChan.Push(events...) } } @@ -570,11 +578,12 @@ func (s *sink) getAllTableNames(ts uint64) []*commonEvent.SchemaTableName { } func (s *sink) Close() { - s.isNormal.Store(false) s.close() + // Close the shared Sarama client before producer facades to help unblock + // their broker operations during shutdown. + s.comp.close() s.ddlProducer.Close() s.dmlProducer.Close() - s.comp.close() s.statistics.Close() } diff --git a/downstreamadapter/sink/kafka/sink_test.go b/downstreamadapter/sink/kafka/sink_test.go index e518a35609..8f82b96b46 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -245,8 +245,10 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { cause := errors.ErrKafkaSendMessage.GenWithStackByArgs() factory.EXPECT().AsyncProducer(gomock.Any()).Return(nil, cause) - adminClient.EXPECT().Close() - topicManager.EXPECT().Close() + gomock.InOrder( + topicManager.EXPECT().Close(), + adminClient.EXPECT().Close(), + ) kafkaSink, err := newWithComponents( t.Context(), @@ -270,9 +272,11 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { factory.EXPECT().AsyncProducer(gomock.Any()).Return(asyncProducer, nil) factory.EXPECT().SyncProducer(gomock.Any()).Return(nil, cause) - asyncProducer.EXPECT().Close() - adminClient.EXPECT().Close() - topicManager.EXPECT().Close() + gomock.InOrder( + topicManager.EXPECT().Close(), + adminClient.EXPECT().Close(), + asyncProducer.EXPECT().Close(), + ) kafkaSink, err := newWithComponents( t.Context(), @@ -298,10 +302,12 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { factory.EXPECT().AsyncProducer(gomock.Any()).Return(asyncProducer, nil) factory.EXPECT().SyncProducer(gomock.Any()).Return(syncProducer, nil) factory.EXPECT().MetricsCollector(adminClient).Return(noopMetricsCollector{}) - asyncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }) - syncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }) - adminClient.EXPECT().Close().Do(func() { closeCount.Add(1) }) - topicManager.EXPECT().Close().Do(func() { closeCount.Add(1) }) + gomock.InOrder( + topicManager.EXPECT().Close().Do(func() { closeCount.Add(1) }), + adminClient.EXPECT().Close().Do(func() { closeCount.Add(1) }), + syncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }), + asyncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }), + ) kafkaSink, err := newWithComponents( t.Context(), @@ -315,14 +321,11 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { require.Zero(t, closeCount.Load()) require.True(t, kafkaSink.IsNormal()) - kafkaSink.isNormal.Store(false) - kafkaSink.AddDMLEvent(&commonEvent.DMLEvent{}) - require.Zero(t, kafkaSink.eventChan.Len()) - kafkaSink.isNormal.Store(true) - kafkaSink.Close() require.Equal(t, int64(4), closeCount.Load()) require.False(t, kafkaSink.IsNormal()) + kafkaSink.AddDMLEvent(&commonEvent.DMLEvent{}) + require.Zero(t, kafkaSink.eventChan.Len()) _, ok, err := kafkaSink.eventChan.GetWithContext(t.Context()) require.NoError(t, err) diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index dfb61ff05f..c72fd24374 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -27,18 +27,11 @@ import ( type saramaAdminClient struct { changefeed common.ChangeFeedID - // client is the underlying sarama client created for this admin wrapper. - // It must be closed to stop background goroutines (e.g. metadata updater) and release memory. - client saramaClient + // client is shared by the runtime admin and producers. admin.Close closes it. + client sarama.Client admin saramaClusterAdmin } -type saramaClient interface { - Brokers() []*sarama.Broker - Partitions(topic string) ([]int32, error) - Close() error -} - type saramaClusterAdmin interface { DescribeCluster() (brokers []*sarama.Broker, controllerID int32, err error) DescribeConfig(resource sarama.ConfigResource) ([]sarama.ConfigEntry, error) @@ -209,24 +202,12 @@ func (a *saramaAdminClient) CreateTopic(detail *TopicDetail) error { } func (a *saramaAdminClient) Close() { - // For admins created via sarama.NewClusterAdminFromClient, admin.Close() takes care - // of closing the underlying client as well. Fall back to closing the client directly - // only when admin is unexpectedly nil. - if a.admin != nil { - if err := a.admin.Close(); err != nil { - log.Warn("kafka admin client close failed", - zap.String("keyspace", a.changefeed.Keyspace()), - zap.String("changefeed", a.changefeed.Name()), - zap.Error(err)) - } - return - } - if a.client != nil { - if err := a.client.Close(); err != nil { - log.Warn("kafka client close failed", - zap.String("keyspace", a.changefeed.Keyspace()), - zap.String("changefeed", a.changefeed.Name()), - zap.Error(err)) - } + // NewClusterAdminFromClient transfers the close responsibility to the admin. + // Closing it also releases the client shared by the runtime producers. + if err := a.admin.Close(); err != nil { + log.Warn("kafka admin client close failed", + zap.String("keyspace", a.changefeed.Keyspace()), + zap.String("changefeed", a.changefeed.Name()), + zap.Error(err)) } } diff --git a/pkg/sink/kafka/sarama_admin_mock.go b/pkg/sink/kafka/sarama_admin_mock.go index 588e08986b..0b5a40bce8 100644 --- a/pkg/sink/kafka/sarama_admin_mock.go +++ b/pkg/sink/kafka/sarama_admin_mock.go @@ -11,72 +11,6 @@ import ( gomock "github.com/golang/mock/gomock" ) -// MocksaramaClient is a mock of saramaClient interface. -type MocksaramaClient struct { - ctrl *gomock.Controller - recorder *MocksaramaClientMockRecorder -} - -// MocksaramaClientMockRecorder is the mock recorder for MocksaramaClient. -type MocksaramaClientMockRecorder struct { - mock *MocksaramaClient -} - -// NewMocksaramaClient creates a new mock instance. -func NewMocksaramaClient(ctrl *gomock.Controller) *MocksaramaClient { - mock := &MocksaramaClient{ctrl: ctrl} - mock.recorder = &MocksaramaClientMockRecorder{mock} - return mock -} - -// EXPECT returns an object that allows the caller to indicate expected use. -func (m *MocksaramaClient) EXPECT() *MocksaramaClientMockRecorder { - return m.recorder -} - -// Brokers mocks base method. -func (m *MocksaramaClient) Brokers() []*sarama.Broker { - m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "Brokers") - ret0, _ := ret[0].([]*sarama.Broker) - return ret0 -} - -// Brokers indicates an expected call of Brokers. -func (mr *MocksaramaClientMockRecorder) Brokers() *gomock.Call { - mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Brokers", reflect.TypeOf((*MocksaramaClient)(nil).Brokers)) -} - -// Close mocks base method. -func (m *MocksaramaClient) Close() error { - m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "Close") - ret0, _ := ret[0].(error) - return ret0 -} - -// Close indicates an expected call of Close. -func (mr *MocksaramaClientMockRecorder) Close() *gomock.Call { - mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Close", reflect.TypeOf((*MocksaramaClient)(nil).Close)) -} - -// Partitions mocks base method. -func (m *MocksaramaClient) Partitions(topic string) ([]int32, error) { - m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "Partitions", topic) - ret0, _ := ret[0].([]int32) - ret1, _ := ret[1].(error) - return ret0, ret1 -} - -// Partitions indicates an expected call of Partitions. -func (mr *MocksaramaClientMockRecorder) Partitions(topic interface{}) *gomock.Call { - mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Partitions", reflect.TypeOf((*MocksaramaClient)(nil).Partitions), topic) -} - // MocksaramaClusterAdmin is a mock of saramaClusterAdmin interface. type MocksaramaClusterAdmin struct { ctrl *gomock.Controller diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index efb34ea361..1a1d6f2754 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -416,43 +416,13 @@ func TestCreateTopic(t *testing.T) { } func TestAdminClientClose(t *testing.T) { - tests := []struct { - name string - setup func(*gomock.Controller) *saramaAdminClient - }{ - { - name: "uses admin close", - setup: func(ctrl *gomock.Controller) *saramaAdminClient { - client := NewMocksaramaClient(ctrl) - admin := NewMocksaramaClusterAdmin(ctrl) - admin.EXPECT().Close().Return(nil) - client.EXPECT().Close().Times(0) - return &saramaAdminClient{ - changefeed: common.NewChangeFeedIDWithName("test", "default"), - client: client, - admin: admin, - } - }, - }, - { - name: "falls back to client when admin is nil", - setup: func(ctrl *gomock.Controller) *saramaAdminClient { - client := NewMocksaramaClient(ctrl) - client.EXPECT().Close().Return(nil) - return &saramaAdminClient{ - changefeed: common.NewChangeFeedIDWithName("test", "default"), - client: client, - } - }, - }, + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().Close().Return(nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, } - for _, test := range tests { - t.Run(test.name, func(t *testing.T) { - ctrl := gomock.NewController(t) - adminClient := test.setup(ctrl) - - require.NotPanics(t, func() { adminClient.Close() }) - }) - } + require.NotPanics(t, func() { client.Close() }) } diff --git a/pkg/sink/kafka/sarama_async_producer.go b/pkg/sink/kafka/sarama_async_producer.go index 5a9ea691e1..16c99d8d9e 100644 --- a/pkg/sink/kafka/sarama_async_producer.go +++ b/pkg/sink/kafka/sarama_async_producer.go @@ -27,7 +27,6 @@ import ( ) type saramaAsyncProducer struct { - client sarama.Client producer sarama.AsyncProducer changefeedID common.ChangeFeedID @@ -47,34 +46,13 @@ func (p *saramaAsyncProducer) Close() { // Safety: // * If the kafka cluster is running well, it will be closed as soon as possible. // Also, we cancel all table pipelines before closed, so it's safe. - // * If there is a problem with the kafka cluster, it will shut down the client first, - // which means no more data will be sent because the connection to the broker is dropped. - // Also, we cancel all table pipelines before closed, so it's safe. + // * The shared client is closed by the admin during sink shutdown, which helps + // unblock broker operations before this producer is closed. // * For Kafka Sink, duplicate data is acceptable. // * There is a risk of goroutine leakage, but it is acceptable and our main // goal is not to get stuck with the processor tick. - // `client` is mainly used by `asyncProducer` to fetch metadata and perform other related - // operations. When we close the `kafkaSaramaProducer`, - // there is no need for TiCDC to make sure that all buffered messages are flushed. - // Consider the situation where the broker is irresponsive. If the client were not - // closed, `asyncProducer.Close()` would waste a mount of time to try flush all messages. - // To prevent the scenario mentioned above, close the client first. start := time.Now() - if err := p.client.Close(); err != nil { - log.Warn("kafka async producer client close failed", - zap.String("keyspace", p.changefeedID.Keyspace()), - zap.String("changefeed", p.changefeedID.Name()), - zap.Duration("duration", time.Since(start)), - zap.Error(err)) - } else { - log.Info("kafka async producer client closed", - zap.String("keyspace", p.changefeedID.Keyspace()), - zap.String("changefeed", p.changefeedID.Name()), - zap.Duration("duration", time.Since(start))) - } - - start = time.Now() if err := p.producer.Close(); err != nil { log.Warn("kafka async producer close failed", zap.String("keyspace", p.changefeedID.Keyspace()), diff --git a/pkg/sink/kafka/sarama_factory.go b/pkg/sink/kafka/sarama_factory.go index e86cdc78a3..e4e1ad5ca7 100644 --- a/pkg/sink/kafka/sarama_factory.go +++ b/pkg/sink/kafka/sarama_factory.go @@ -30,6 +30,7 @@ type saramaFactory struct { changefeedID common.ChangeFeedID option *options metricRegistry metrics.Registry + client sarama.Client } // NewSaramaFactory constructs a Factory with sarama implementation. @@ -83,7 +84,7 @@ func NewSaramaFactory( }, nil } -func newAdminClient(changefeedID common.ChangeFeedID, endpoints []string, config *sarama.Config) (AdminClient, error) { +func newAdminClient(changefeedID common.ChangeFeedID, endpoints []string, config *sarama.Config) (*saramaAdminClient, error) { start := time.Now() client, err := sarama.NewClient(endpoints, config) duration := time.Since(start) @@ -107,8 +108,7 @@ func newAdminClient(changefeedID common.ChangeFeedID, endpoints []string, config zap.Duration("duration", duration)) } if err != nil { - // `sarama.NewClusterAdminFromClient` does not take ownership of the client, - // so we need to close it on failures to avoid leaking background goroutines. + // No admin exists to close the client when construction fails. _ = client.Close() return nil, errors.WrapError(errors.ErrNewKafkaSink, err) } @@ -124,32 +124,26 @@ func (f *saramaFactory) AdminClient(ctx context.Context) (AdminClient, error) { if err != nil { return nil, err } - return newAdminClient(f.changefeedID, f.option.BrokerEndpoints, config) -} + config.MetricRegistry = f.metricRegistry -// SyncProducer returns a Sync SyncProducer, -// it should be the caller's responsibility to close the producer -func (f *saramaFactory) SyncProducer(ctx context.Context) (SyncProducer, error) { - config, err := newSaramaConfig(ctx, f.option) + admin, err := newAdminClient(f.changefeedID, f.option.BrokerEndpoints, config) if err != nil { return nil, err } - config.MetricRegistry = f.metricRegistry + f.client = admin.client + return admin, nil +} - client, err := sarama.NewClient(f.option.BrokerEndpoints, config) +// SyncProducer returns a Sync SyncProducer, +// it should be the caller's responsibility to close the producer +func (f *saramaFactory) SyncProducer(context.Context) (SyncProducer, error) { + p, err := sarama.NewSyncProducerFromClient(f.client) if err != nil { return nil, errors.WrapError(errors.ErrNewKafkaSink, err) } - p, err := sarama.NewSyncProducerFromClient(client) - if err != nil { - _ = client.Close() - return nil, errors.WrapError(errors.ErrNewKafkaSink, err) - } - return &saramaSyncProducer{ id: f.changefeedID, - client: client, producer: p, closed: atomic.NewBool(false), }, nil @@ -157,25 +151,12 @@ func (f *saramaFactory) SyncProducer(ctx context.Context) (SyncProducer, error) // AsyncProducer return an Async SyncProducer, // it should be the caller's responsibility to close the producer -func (f *saramaFactory) AsyncProducer(ctx context.Context) (AsyncProducer, error) { - config, err := newSaramaConfig(ctx, f.option) - if err != nil { - return nil, err - } - config.MetricRegistry = f.metricRegistry - - client, err := sarama.NewClient(f.option.BrokerEndpoints, config) - if err != nil { - return nil, errors.WrapError(errors.ErrNewKafkaSink, err) - } - - p, err := sarama.NewAsyncProducerFromClient(client) +func (f *saramaFactory) AsyncProducer(context.Context) (AsyncProducer, error) { + p, err := sarama.NewAsyncProducerFromClient(f.client) if err != nil { - _ = client.Close() return nil, errors.WrapError(errors.ErrNewKafkaSink, err) } return &saramaAsyncProducer{ - client: client, producer: p, changefeedID: f.changefeedID, closed: atomic.NewBool(false), diff --git a/pkg/sink/kafka/sarama_sync_producer.go b/pkg/sink/kafka/sarama_sync_producer.go index fcf1c9c258..7c90004705 100644 --- a/pkg/sink/kafka/sarama_sync_producer.go +++ b/pkg/sink/kafka/sarama_sync_producer.go @@ -25,11 +25,6 @@ import ( "go.uber.org/zap" ) -type saramaSyncClient interface { - Brokers() []*sarama.Broker - Close() error -} - type saramaSyncProducerClient interface { SendMessage(msg *sarama.ProducerMessage) (partition int32, offset int64, err error) SendMessages(msgs []*sarama.ProducerMessage) error @@ -38,7 +33,6 @@ type saramaSyncProducerClient interface { type saramaSyncProducer struct { id common.ChangeFeedID - client saramaSyncClient producer saramaSyncProducerClient closed *atomic.Bool } @@ -102,17 +96,6 @@ func (p *saramaSyncProducer) Close() { p.closed.Store(true) start := time.Now() - // sarama.NewSyncProducerFromClient wraps the provided client with a nopCloserClient, - // so producer.Close() alone won't release the underlying client resources. - if p.client != nil { - if err := p.client.Close(); err != nil { - log.Warn("kafka ddl producer client close failed", - zap.String("keyspace", p.id.Keyspace()), - zap.String("changefeed", p.id.Name()), - zap.Duration("duration", time.Since(start)), - zap.Error(err)) - } - } if p.producer != nil { if err := p.producer.Close(); err != nil { log.Error("kafka ddl producer close failed", diff --git a/pkg/sink/kafka/sarama_sync_producer_mock.go b/pkg/sink/kafka/sarama_sync_producer_mock.go index 78671e02f2..4ba58240de 100644 --- a/pkg/sink/kafka/sarama_sync_producer_mock.go +++ b/pkg/sink/kafka/sarama_sync_producer_mock.go @@ -11,57 +11,6 @@ import ( gomock "github.com/golang/mock/gomock" ) -// MocksaramaSyncClient is a mock of saramaSyncClient interface. -type MocksaramaSyncClient struct { - ctrl *gomock.Controller - recorder *MocksaramaSyncClientMockRecorder -} - -// MocksaramaSyncClientMockRecorder is the mock recorder for MocksaramaSyncClient. -type MocksaramaSyncClientMockRecorder struct { - mock *MocksaramaSyncClient -} - -// NewMocksaramaSyncClient creates a new mock instance. -func NewMocksaramaSyncClient(ctrl *gomock.Controller) *MocksaramaSyncClient { - mock := &MocksaramaSyncClient{ctrl: ctrl} - mock.recorder = &MocksaramaSyncClientMockRecorder{mock} - return mock -} - -// EXPECT returns an object that allows the caller to indicate expected use. -func (m *MocksaramaSyncClient) EXPECT() *MocksaramaSyncClientMockRecorder { - return m.recorder -} - -// Brokers mocks base method. -func (m *MocksaramaSyncClient) Brokers() []*sarama.Broker { - m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "Brokers") - ret0, _ := ret[0].([]*sarama.Broker) - return ret0 -} - -// Brokers indicates an expected call of Brokers. -func (mr *MocksaramaSyncClientMockRecorder) Brokers() *gomock.Call { - mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Brokers", reflect.TypeOf((*MocksaramaSyncClient)(nil).Brokers)) -} - -// Close mocks base method. -func (m *MocksaramaSyncClient) Close() error { - m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "Close") - ret0, _ := ret[0].(error) - return ret0 -} - -// Close indicates an expected call of Close. -func (mr *MocksaramaSyncClientMockRecorder) Close() *gomock.Call { - mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Close", reflect.TypeOf((*MocksaramaSyncClient)(nil).Close)) -} - // MocksaramaSyncProducerClient is a mock of saramaSyncProducerClient interface. type MocksaramaSyncProducerClient struct { ctrl *gomock.Controller diff --git a/pkg/sink/kafka/sarama_sync_producer_test.go b/pkg/sink/kafka/sarama_sync_producer_test.go index e00e0dafa8..41fec47afe 100644 --- a/pkg/sink/kafka/sarama_sync_producer_test.go +++ b/pkg/sink/kafka/sarama_sync_producer_test.go @@ -40,39 +40,16 @@ func TestProducerRejectsSendAfterClose(t *testing.T) { } func TestSyncProducerClose(t *testing.T) { - tests := []struct { - name string - clientCloseErr error - }{ - { - name: "closes client and producer", - }, - { - name: "still closes producer when client close fails", - clientCloseErr: io.ErrClosedPipe, - }, + ctrl := gomock.NewController(t) + producer := NewMocksaramaSyncProducerClient(ctrl) + producer.EXPECT().Close().Return(nil) + p := &saramaSyncProducer{ + id: common.NewChangeFeedIDWithName("test", "default"), + producer: producer, + closed: atomic.NewBool(false), } - for _, test := range tests { - t.Run(test.name, func(t *testing.T) { - ctrl := gomock.NewController(t) - client := NewMocksaramaSyncClient(ctrl) - producer := NewMocksaramaSyncProducerClient(ctrl) - gomock.InOrder( - client.EXPECT().Close().Return(test.clientCloseErr), - producer.EXPECT().Close().Return(nil), - ) - - p := &saramaSyncProducer{ - id: common.NewChangeFeedIDWithName("test", "default"), - client: client, - producer: producer, - closed: atomic.NewBool(false), - } - - p.Close() - }) - } + p.Close() } func TestSyncProducerErrorWrappedOnce(t *testing.T) { From 14ba283830b1b43c8eba5860b267f51f0c2d3476 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Wed, 26 Aug 2026 15:15:45 +0800 Subject: [PATCH 11/13] topic manager add wait group --- .../sink/topicmanager/kafka_topic_manager.go | 14 +++++-- .../topicmanager/kafka_topic_manager_test.go | 39 +++++++++++++++++++ pkg/sink/kafka/admin.go | 26 ++++++++++--- pkg/sink/kafka/sarama_admin_test.go | 28 ++++++++++++- 4 files changed, 97 insertions(+), 10 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 5dc68119e4..21c7ad2fdd 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -46,6 +46,7 @@ type kafkaTopicManager struct { topics sync.Map // cancel is used to cancel the background goroutine. cancel context.CancelFunc + wg sync.WaitGroup } // newKafkaTopicManager creates a topic manager without starting background work. @@ -91,7 +92,11 @@ func GetTopicManagerAndTryCreateTopic( } ctx, cancel := context.WithCancel(ctx) topicManager.cancel = cancel - go topicManager.backgroundRefreshMeta(ctx) + topicManager.wg.Add(1) + go func() { + defer topicManager.wg.Done() + topicManager.backgroundRefreshMeta(ctx) + }() return topicManager, nil } @@ -350,7 +355,10 @@ func (m *kafkaTopicManager) useConfiguredPartitionNum(topicName string, cause er return m.cfg.PartitionNum } -// Close exits the background goroutine. +// Close cancels the background goroutine and waits for it to exit. func (m *kafkaTopicManager) Close() { - m.cancel() + if m.cancel != nil { + m.cancel() + } + m.wg.Wait() } diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 0976339f61..f46c0575e7 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -17,6 +17,7 @@ import ( "context" "io" "testing" + "time" "github.com/IBM/sarama" "github.com/golang/mock/gomock" @@ -260,6 +261,44 @@ func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { require.NotNil(t, manager.(*kafkaTopicManager).cancel) } +func TestKafkaTopicManagerCloseWaitsForBackgroundWork(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + manager := &kafkaTopicManager{cancel: cancel} + canceled := make(chan struct{}) + release := make(chan struct{}) + manager.wg.Add(1) + go func() { + defer manager.wg.Done() + <-ctx.Done() + close(canceled) + <-release + }() + + closed := make(chan struct{}) + go func() { + manager.Close() + close(closed) + }() + + select { + case <-canceled: + case <-time.After(time.Second): + require.FailNow(t, "background work was not canceled") + } + select { + case <-closed: + require.FailNow(t, "Close returned before background work exited") + default: + } + + close(release) + select { + case <-closed: + case <-time.After(time.Second): + require.FailNow(t, "Close did not wait for background work to exit") + } +} + func TestCreateTopicWithTopicDescribeDenied(t *testing.T) { t.Parallel() diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index c72fd24374..712b46fcfa 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -32,6 +32,9 @@ type saramaAdminClient struct { admin saramaClusterAdmin } +// saramaClusterAdmin is the subset of sarama.ClusterAdmin used by TiCDC. +// The narrow interface also lets cleanup tests verify that closing the admin +// releases the client passed to sarama.NewClusterAdminFromClient. type saramaClusterAdmin interface { DescribeCluster() (brokers []*sarama.Broker, controllerID int32, err error) DescribeConfig(resource sarama.ConfigResource) ([]sarama.ConfigEntry, error) @@ -204,10 +207,23 @@ func (a *saramaAdminClient) CreateTopic(detail *TopicDetail) error { func (a *saramaAdminClient) Close() { // NewClusterAdminFromClient transfers the close responsibility to the admin. // Closing it also releases the client shared by the runtime producers. - if err := a.admin.Close(); err != nil { - log.Warn("kafka admin client close failed", - zap.String("keyspace", a.changefeed.Keyspace()), - zap.String("changefeed", a.changefeed.Name()), - zap.Error(err)) + if a.admin != nil { + if err := a.admin.Close(); err != nil { + log.Warn("kafka admin client close failed", + zap.String("keyspace", a.changefeed.Keyspace()), + zap.String("changefeed", a.changefeed.Name()), + zap.Error(err)) + } + return + } + // Preserve cleanup for a partially initialized wrapper. This is also the + // resource-leak regression guard from PR #4437. + if a.client != nil { + if err := a.client.Close(); err != nil { + log.Warn("kafka client close failed", + zap.String("keyspace", a.changefeed.Keyspace()), + zap.String("changefeed", a.changefeed.Name()), + zap.Error(err)) + } } } diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index 1a1d6f2754..d0d6cbd023 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -415,14 +415,38 @@ func TestCreateTopic(t *testing.T) { } } -func TestAdminClientClose(t *testing.T) { +type closeTrackingSaramaClient struct { + sarama.Client + closed bool +} + +func (c *closeTrackingSaramaClient) Close() error { + c.closed = true + return nil +} + +func TestSaramaAdminClientCloseDelegatesClientCleanupToAdmin(t *testing.T) { ctrl := gomock.NewController(t) + underlyingClient := &closeTrackingSaramaClient{} admin := NewMocksaramaClusterAdmin(ctrl) - admin.EXPECT().Close().Return(nil) + admin.EXPECT().Close().DoAndReturn(underlyingClient.Close) client := &saramaAdminClient{ changefeed: common.NewChangeFeedIDWithName("test", "default"), + client: underlyingClient, admin: admin, } + client.Close() + require.True(t, underlyingClient.closed) +} + +func TestSaramaAdminClientCloseFallsBackToClientWhenAdminIsNil(t *testing.T) { + underlyingClient := &closeTrackingSaramaClient{} + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + client: underlyingClient, + } + require.NotPanics(t, func() { client.Close() }) + require.True(t, underlyingClient.closed) } From 636c3fe07561d84d07cead963a4ee508518c1d31 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Wed, 26 Aug 2026 18:13:27 +0800 Subject: [PATCH 12/13] kafka: refine client ownership and retry handling --- downstreamadapter/sink/kafka/helper.go | 4 +- downstreamadapter/sink/kafka/sink.go | 10 +-- downstreamadapter/sink/kafka/sink_test.go | 9 ++- .../sink/topicmanager/kafka_topic_manager.go | 2 +- .../topicmanager/kafka_topic_manager_test.go | 23 +----- pkg/sink/kafka/admin.go | 16 ++-- pkg/sink/kafka/factory.go | 6 +- pkg/sink/kafka/factory_mock.go | 12 +++ pkg/sink/kafka/sarama_admin_test.go | 44 +++++++---- pkg/sink/kafka/sarama_async_producer.go | 2 +- pkg/sink/kafka/sarama_factory.go | 75 ++++++++++--------- 11 files changed, 112 insertions(+), 91 deletions(-) diff --git a/downstreamadapter/sink/kafka/helper.go b/downstreamadapter/sink/kafka/helper.go index b3353afe25..2dfe3e3430 100644 --- a/downstreamadapter/sink/kafka/helper.go +++ b/downstreamadapter/sink/kafka/helper.go @@ -45,8 +45,8 @@ func (c components) close() { if c.topicManager != nil { c.topicManager.Close() } - if c.adminClient != nil { - c.adminClient.Close() + if c.factory != nil { + c.factory.Close() } if c.claimCheck != nil { c.claimCheck.Close() diff --git a/downstreamadapter/sink/kafka/sink.go b/downstreamadapter/sink/kafka/sink.go index a156f38021..f34d3971a8 100644 --- a/downstreamadapter/sink/kafka/sink.go +++ b/downstreamadapter/sink/kafka/sink.go @@ -122,12 +122,12 @@ func Verify(ctx context.Context, changefeedID common.ChangeFeedID, uri *url.URL, if err != nil { return err } + defer factory.Close() adminClient, err := factory.AdminClient(ctx) if err != nil { return err } - defer adminClient.Close() err = topicmanager.EnsureTopic(ctx, changefeedID, topic, options.DeriveTopicConfig(), adminClient) if err != nil { @@ -168,8 +168,8 @@ func newWithComponents( if err == nil { return } - // Closing the admin releases the shared Sarama client and unblocks - // producer shutdown when Kafka is unhealthy. + // Release shared Kafka resources first to help unblock producer shutdown + // when Kafka is unhealthy. comp.close() if syncProducer != nil { syncProducer.Close() @@ -579,8 +579,8 @@ func (s *sink) getAllTableNames(ts uint64) []*commonEvent.SchemaTableName { func (s *sink) Close() { s.close() - // Close the shared Sarama client before producer facades to help unblock - // their broker operations during shutdown. + // Release shared Kafka resources before closing producers to help unblock + // their shutdown. s.comp.close() s.ddlProducer.Close() s.dmlProducer.Close() diff --git a/downstreamadapter/sink/kafka/sink_test.go b/downstreamadapter/sink/kafka/sink_test.go index 8f82b96b46..e9b77c342a 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -112,7 +112,7 @@ func TestVerifyInvalidConfig(t *testing.T) { factory.EXPECT().AdminClient(gomock.Any()).Return(adminClient, nil), adminClient.EXPECT().GetTopicsMeta([]string{kafkaSinkTestTopic}, true).Return( map[string]kafka.TopicDetail{kafkaSinkTestTopic: {Name: kafkaSinkTestTopic}}, nil), - adminClient.EXPECT().Close(), + factory.EXPECT().Close(), ) originalCreateKafkaFactory := createKafkaFactory @@ -247,7 +247,7 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { factory.EXPECT().AsyncProducer(gomock.Any()).Return(nil, cause) gomock.InOrder( topicManager.EXPECT().Close(), - adminClient.EXPECT().Close(), + factory.EXPECT().Close(), ) kafkaSink, err := newWithComponents( @@ -274,7 +274,7 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { factory.EXPECT().SyncProducer(gomock.Any()).Return(nil, cause) gomock.InOrder( topicManager.EXPECT().Close(), - adminClient.EXPECT().Close(), + factory.EXPECT().Close(), asyncProducer.EXPECT().Close(), ) @@ -304,7 +304,7 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { factory.EXPECT().MetricsCollector(adminClient).Return(noopMetricsCollector{}) gomock.InOrder( topicManager.EXPECT().Close().Do(func() { closeCount.Add(1) }), - adminClient.EXPECT().Close().Do(func() { closeCount.Add(1) }), + factory.EXPECT().Close().Do(func() { closeCount.Add(1) }), syncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }), asyncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }), ) @@ -607,6 +607,7 @@ func newKafkaSinkForTest( factory.EXPECT().AsyncProducer(gomock.Any()).Return(asyncProducer, nil) factory.EXPECT().SyncProducer(gomock.Any()).Return(syncProducer, nil) factory.EXPECT().MetricsCollector(nil).Return(noopMetricsCollector{}) + factory.EXPECT().Close().AnyTimes() kafkaSink, err := newWithComponents(ctx, changefeedID, common.DefaultKeyspaceID, protocol, components{ encoderGroup: encoderGroup, diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 21c7ad2fdd..757e927a7e 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -220,7 +220,7 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( retry.WithBackoffMaxDelay(1000), retry.WithMaxTries(6), retry.WithIsRetryableErr(func(err error) bool { - return !kafka.IsUnretryableTopicMetadataError(err) + return !kafka.IsUnretryableKafkaError(err) }), ) if err != nil { diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index f46c0575e7..2c351425c8 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -75,23 +75,13 @@ func TestCreateTopic(t *testing.T) { adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).DoAndReturn( func([]string, bool) (map[string]kafka.TopicDetail, error) { if createdTopic == nil { - return nil, errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrUnknownTopicOrPartition, - "describe-topic", - "new-topic", - ) + return nil, errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ErrUnknownTopicOrPartition, "describe-topic", "new-topic") } postCreateDescribeCount++ _, cached := manager.topics.Load("new-topic") require.False(t, cached) if postCreateDescribeCount == 1 { - return nil, errors.WrapError( - errors.ErrKafkaAdminAPI, - io.EOF, - "describe-topic", - "new-topic", - ) + return nil, errors.WrapError(errors.ErrKafkaAdminAPI, io.EOF, "describe-topic", "new-topic") } return map[string]kafka.TopicDetail{ createdTopic.Name: topicDetail(createdTopic.Name, createdTopic.NumPartitions), @@ -211,19 +201,14 @@ func TestCreateTopicValidatesReplicationFactor(t *testing.T) { require.ErrorContains(t, err, "`replication-factor` 1 is smaller than the `min.insync.replicas` 2 of broker") } -func TestWaitUntilTopicVisibleStopsOnNonRetryableError(t *testing.T) { +func TestWaitUntilTopicVisibleUnretryableError(t *testing.T) { t.Parallel() ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{"invalid-topic"}, false).Return( nil, - errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrInvalidTopic, - "describe-topic", - "invalid-topic", - ), + errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ErrInvalidTopic, "describe-topic", "invalid-topic"), ).Times(1) manager := newKafkaTopicManager( "invalid-topic", diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index 712b46fcfa..311146f1eb 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -27,7 +27,8 @@ import ( type saramaAdminClient struct { changefeed common.ChangeFeedID - // client is shared by the runtime admin and producers. admin.Close closes it. + // client is set for Admin methods that use the client shared by the Factory. + // Configuration probing only uses admin and leaves client nil. client sarama.Client admin saramaClusterAdmin } @@ -152,9 +153,10 @@ func IsAuthorizationFailed(err error) bool { errors.Is(err, sarama.ErrClusterAuthorizationFailed) } -// IsUnretryableTopicMetadataError reports whether a Kafka metadata request -// requires a configuration, credential, permission, or request change to succeed. -func IsUnretryableTopicMetadataError(err error) bool { +// IsUnretryableKafkaError reports whether err is not retryable. +// See Apache Kafka protocol error definitions: +// https://kafka.apache.org/38/generated/protocol_errors.html +func IsUnretryableKafkaError(err error) bool { if IsAuthorizationFailed(err) || errors.Is(err, errors.ErrKafkaInvalidConfig) || errors.Is(err, sarama.ErrInvalidTopic) || @@ -205,8 +207,8 @@ func (a *saramaAdminClient) CreateTopic(detail *TopicDetail) error { } func (a *saramaAdminClient) Close() { - // NewClusterAdminFromClient transfers the close responsibility to the admin. - // Closing it also releases the client shared by the runtime producers. + // Sarama closes the client passed to NewClusterAdminFromClient when the + // cluster admin is closed. if a.admin != nil { if err := a.admin.Close(); err != nil { log.Warn("kafka admin client close failed", @@ -216,7 +218,7 @@ func (a *saramaAdminClient) Close() { } return } - // Preserve cleanup for a partially initialized wrapper. This is also the + // Preserve cleanup for a partially initialized admin. This is also the // resource-leak regression guard from PR #4437. if a.client != nil { if err := a.client.Close(); err != nil { diff --git a/pkg/sink/kafka/factory.go b/pkg/sink/kafka/factory.go index c19089de4c..de332c492d 100644 --- a/pkg/sink/kafka/factory.go +++ b/pkg/sink/kafka/factory.go @@ -29,6 +29,8 @@ type Factory interface { AsyncProducer(ctx context.Context) (AsyncProducer, error) // MetricsCollector returns the kafka metrics collector MetricsCollector(adminClient AdminClient) MetricsCollector + // Close releases resources shared by all components created by this factory. + Close() } // SyncProducer is the kafka sync producer @@ -46,7 +48,6 @@ type SyncProducer interface { // Close shuts down the producer; you must call this function before a producer // object passes out of scope, as it may otherwise leak memory. - // You must call this before calling Close on the underlying client. Close() } @@ -55,8 +56,7 @@ type AsyncProducer interface { // Close shuts down the producer and waits for any buffered messages to be // flushed. You must call this function before a producer object passes out of // scope, as it may otherwise leak memory. You must call this before process - // shutting down, or you may lose messages. You must call this before calling - // Close on the underlying client. + // shutting down, or you may lose messages. Close() // AsyncSend is the input channel for the user to write messages to that they diff --git a/pkg/sink/kafka/factory_mock.go b/pkg/sink/kafka/factory_mock.go index ecc8fe131c..fabd1fc208 100644 --- a/pkg/sink/kafka/factory_mock.go +++ b/pkg/sink/kafka/factory_mock.go @@ -65,6 +65,18 @@ func (mr *MockFactoryMockRecorder) AsyncProducer(ctx interface{}) *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AsyncProducer", reflect.TypeOf((*MockFactory)(nil).AsyncProducer), ctx) } +// Close mocks base method. +func (m *MockFactory) Close() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "Close") +} + +// Close indicates an expected call of Close. +func (mr *MockFactoryMockRecorder) Close() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Close", reflect.TypeOf((*MockFactory)(nil).Close)) +} + // MetricsCollector mocks base method. func (m *MockFactory) MetricsCollector(adminClient AdminClient) MetricsCollector { m.ctrl.T.Helper() diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index d0d6cbd023..e626287b5d 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -312,7 +312,7 @@ func TestIsAuthorizationFailed(t *testing.T) { } } -func TestIsUnretryableTopicMetadataError(t *testing.T) { +func TestIsUnretryableKafkaError(t *testing.T) { t.Parallel() tests := []struct { @@ -330,12 +330,7 @@ func TestIsUnretryableTopicMetadataError(t *testing.T) { {name: "unknown broker error", err: sarama.ErrUnknown}, { name: "wrapped unknown topic", - err: errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrUnknownTopicOrPartition, - "describe-topic", - "test-topic", - ), + err: errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ErrUnknownTopicOrPartition, "describe-topic", "test-topic"), }, {name: "context cancellation", err: context.Canceled}, {name: "TiCDC invalid config", err: errors.ErrKafkaInvalidConfig.GenWithStack("invalid config"), unretryable: true}, @@ -350,20 +345,15 @@ func TestIsUnretryableTopicMetadataError(t *testing.T) { {name: "invalid request", err: sarama.ErrInvalidRequest, unretryable: true}, {name: "client configuration error", err: sarama.ConfigurationError("invalid client config"), unretryable: true}, { - name: "wrapped invalid topic", - err: errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrInvalidTopic, - "describe-topic", - "test-topic", - ), + name: "wrapped invalid topic", + err: errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ErrInvalidTopic, "describe-topic", "test-topic"), unretryable: true, }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { - require.Equal(t, test.unretryable, IsUnretryableTopicMetadataError(test.err)) + require.Equal(t, test.unretryable, IsUnretryableKafkaError(test.err)) }) } } @@ -425,6 +415,14 @@ func (c *closeTrackingSaramaClient) Close() error { return nil } +func (c *closeTrackingSaramaClient) Controller() (*sarama.Broker, error) { + return &sarama.Broker{}, nil +} + +func (c *closeTrackingSaramaClient) Config() *sarama.Config { + return sarama.NewConfig() +} + func TestSaramaAdminClientCloseDelegatesClientCleanupToAdmin(t *testing.T) { ctrl := gomock.NewController(t) underlyingClient := &closeTrackingSaramaClient{} @@ -450,3 +448,19 @@ func TestSaramaAdminClientCloseFallsBackToClientWhenAdminIsNil(t *testing.T) { require.NotPanics(t, func() { client.Close() }) require.True(t, underlyingClient.closed) } + +func TestSaramaFactoryOwnsSharedClient(t *testing.T) { + underlyingClient := &closeTrackingSaramaClient{} + factory := &saramaFactory{ + changefeedID: common.NewChangeFeedIDWithName("test", "default"), + client: underlyingClient, + } + + adminClient, err := factory.AdminClient(t.Context()) + require.NoError(t, err) + require.Same(t, underlyingClient, adminClient.(*saramaAdminClient).client) + require.False(t, underlyingClient.closed) + + factory.Close() + require.True(t, underlyingClient.closed) +} diff --git a/pkg/sink/kafka/sarama_async_producer.go b/pkg/sink/kafka/sarama_async_producer.go index 16c99d8d9e..1371b21e17 100644 --- a/pkg/sink/kafka/sarama_async_producer.go +++ b/pkg/sink/kafka/sarama_async_producer.go @@ -46,7 +46,7 @@ func (p *saramaAsyncProducer) Close() { // Safety: // * If the kafka cluster is running well, it will be closed as soon as possible. // Also, we cancel all table pipelines before closed, so it's safe. - // * The shared client is closed by the admin during sink shutdown, which helps + // * The shared client is closed by the Factory during sink shutdown, which helps // unblock broker operations before this producer is closed. // * For Kafka Sink, duplicate data is acceptable. // * There is a risk of goroutine leakage, but it is acceptable and our main diff --git a/pkg/sink/kafka/sarama_factory.go b/pkg/sink/kafka/sarama_factory.go index e4e1ad5ca7..dfa4f09003 100644 --- a/pkg/sink/kafka/sarama_factory.go +++ b/pkg/sink/kafka/sarama_factory.go @@ -28,7 +28,6 @@ import ( type saramaFactory struct { changefeedID common.ChangeFeedID - option *options metricRegistry metrics.Registry client sarama.Client } @@ -52,10 +51,19 @@ func NewSaramaFactory( return nil, err } - admin, err := newAdminClient(changefeedID, o.BrokerEndpoints, config) + start = time.Now() + clusterAdmin, err := sarama.NewClusterAdmin(o.BrokerEndpoints, config) + duration = time.Since(start) + if duration > 2*time.Second { + log.Warn("kafka admin client initialization is slow", + zap.String("keyspace", changefeedID.Keyspace()), + zap.String("changefeed", changefeedID.Name()), + zap.Duration("duration", duration)) + } if err != nil { - return nil, err + return nil, errors.WrapError(errors.ErrNewKafkaSink, err) } + admin := &saramaAdminClient{changefeed: changefeedID, admin: clusterAdmin} defer func() { admin.Close() }() @@ -77,17 +85,14 @@ func NewSaramaFactory( zap.Duration("readTimeout", o.ReadTimeout), zap.Duration("writeTimeout", o.WriteTimeout)) - return &saramaFactory{ - changefeedID: changefeedID, - option: o, - metricRegistry: metrics.NewRegistry(), - }, nil -} + config, err = newSaramaConfig(ctx, o) + if err != nil { + return nil, err + } -func newAdminClient(changefeedID common.ChangeFeedID, endpoints []string, config *sarama.Config) (*saramaAdminClient, error) { - start := time.Now() - client, err := sarama.NewClient(endpoints, config) - duration := time.Since(start) + start = time.Now() + client, err := sarama.NewClient(o.BrokerEndpoints, config) + duration = time.Since(start) if duration > 2*time.Second { log.Warn("kafka client initialization is slow", zap.String("keyspace", changefeedID.Keyspace()), @@ -98,40 +103,42 @@ func newAdminClient(changefeedID common.ChangeFeedID, endpoints []string, config return nil, errors.WrapError(errors.ErrNewKafkaSink, err) } - start = time.Now() - admin, err := sarama.NewClusterAdminFromClient(client) - duration = time.Since(start) + return &saramaFactory{ + changefeedID: changefeedID, + metricRegistry: config.MetricRegistry, + client: client, + }, nil +} + +func (f *saramaFactory) AdminClient(context.Context) (AdminClient, error) { + start := time.Now() + admin, err := sarama.NewClusterAdminFromClient(f.client) + duration := time.Since(start) if duration > 2*time.Second { log.Warn("kafka admin client initialization is slow", - zap.String("keyspace", changefeedID.Keyspace()), - zap.String("changefeed", changefeedID.Name()), + zap.String("keyspace", f.changefeedID.Keyspace()), + zap.String("changefeed", f.changefeedID.Name()), zap.Duration("duration", duration)) } if err != nil { - // No admin exists to close the client when construction fails. - _ = client.Close() return nil, errors.WrapError(errors.ErrNewKafkaSink, err) } return &saramaAdminClient{ - client: client, + client: f.client, admin: admin, - changefeed: changefeedID, + changefeed: f.changefeedID, }, nil } -func (f *saramaFactory) AdminClient(ctx context.Context) (AdminClient, error) { - config, err := newSaramaConfig(ctx, f.option) - if err != nil { - return nil, err - } - config.MetricRegistry = f.metricRegistry - - admin, err := newAdminClient(f.changefeedID, f.option.BrokerEndpoints, config) - if err != nil { - return nil, err +func (f *saramaFactory) Close() { + if f.client != nil { + if err := f.client.Close(); err != nil { + log.Warn("kafka client close failed", + zap.String("keyspace", f.changefeedID.Keyspace()), + zap.String("changefeed", f.changefeedID.Name()), + zap.Error(err)) + } } - f.client = admin.client - return admin, nil } // SyncProducer returns a Sync SyncProducer, From b745ad6f87c416ccb49614eea9430f4bbb3cef47 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Wed, 26 Aug 2026 19:01:19 +0800 Subject: [PATCH 13/13] kafka: refine recovery tests and goroutine startup --- downstreamadapter/sink/kafka/sink_test.go | 17 +++++++++++++++++ .../sink/topicmanager/kafka_topic_manager.go | 6 ++---- .../topicmanager/kafka_topic_manager_test.go | 19 +++++-------------- pkg/sink/kafka/sarama_admin_test.go | 1 + 4 files changed, 25 insertions(+), 18 deletions(-) diff --git a/downstreamadapter/sink/kafka/sink_test.go b/downstreamadapter/sink/kafka/sink_test.go index e9b77c342a..13a6073bdc 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -412,6 +412,23 @@ func TestKafkaSinkDML(t *testing.T) { require.Equal(t, cause, kafkaSink.calculateKeyPartitions(t.Context())) }) + + t.Run("canceled after topic lookup", func(t *testing.T) { + dmlEvent := eventHelper.DML2Event("test", "t", "insert into t values (4, 'four')") + ctx, cancel := context.WithCancelCause(t.Context()) + kafkaSink, topicManager, _, _ := newKafkaSinkForTest( + t, ctx, config.ProtocolOpen, &config.SinkConfig{}) + cause := errors.ErrKafkaSinkClosed.GenWithStackByArgs() + topicManager.EXPECT().GetPartitionNum(gomock.Any(), kafkaSinkTestTopic). + DoAndReturn(func(context.Context, string) (int32, error) { + cancel(cause) + return 1, nil + }) + kafkaSink.AddDMLEvent(dmlEvent) + + require.Equal(t, cause, kafkaSink.calculateKeyPartitions(ctx)) + require.Zero(t, kafkaSink.rowChan.Len()) + }) } func TestKafkaSinkDDL(t *testing.T) { diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 757e927a7e..5d224831a6 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -92,11 +92,9 @@ func GetTopicManagerAndTryCreateTopic( } ctx, cancel := context.WithCancel(ctx) topicManager.cancel = cancel - topicManager.wg.Add(1) - go func() { - defer topicManager.wg.Done() + topicManager.wg.Go(func() { topicManager.backgroundRefreshMeta(ctx) - }() + }) return topicManager, nil } diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 2c351425c8..230c7fa767 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -29,13 +29,6 @@ import ( const kafkaTopicManagerTestTopic = "mock_topic" -func topicDetail(topic string, partitionNum int32) kafka.TopicDetail { - return kafka.TopicDetail{ - Name: topic, - NumPartitions: partitionNum, - } -} - func TestCreateTopic(t *testing.T) { t.Parallel() @@ -48,7 +41,7 @@ func TestCreateTopic(t *testing.T) { adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{kafkaTopicManagerTestTopic}, true).Return( map[string]kafka.TopicDetail{ - kafkaTopicManagerTestTopic: topicDetail(kafkaTopicManagerTestTopic, 2), + kafkaTopicManagerTestTopic: {Name: kafkaTopicManagerTestTopic, NumPartitions: 2}, }, nil) manager := newKafkaTopicManager( kafkaTopicManagerTestTopic, @@ -84,7 +77,7 @@ func TestCreateTopic(t *testing.T) { return nil, errors.WrapError(errors.ErrKafkaAdminAPI, io.EOF, "describe-topic", "new-topic") } return map[string]kafka.TopicDetail{ - createdTopic.Name: topicDetail(createdTopic.Name, createdTopic.NumPartitions), + createdTopic.Name: {Name: createdTopic.Name, NumPartitions: createdTopic.NumPartitions}, }, nil }).Times(3) adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( @@ -230,7 +223,7 @@ func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{"existing-topic"}, true).Return( map[string]kafka.TopicDetail{ - "existing-topic": topicDetail("existing-topic", 2), + "existing-topic": {Name: "existing-topic", NumPartitions: 2}, }, nil) manager, err := GetTopicManagerAndTryCreateTopic( @@ -251,13 +244,11 @@ func TestKafkaTopicManagerCloseWaitsForBackgroundWork(t *testing.T) { manager := &kafkaTopicManager{cancel: cancel} canceled := make(chan struct{}) release := make(chan struct{}) - manager.wg.Add(1) - go func() { - defer manager.wg.Done() + manager.wg.Go(func() { <-ctx.Done() close(canceled) <-release - }() + }) closed := make(chan struct{}) go func() { diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index e626287b5d..6d11863fca 100644 --- a/pkg/sink/kafka/sarama_admin_test.go +++ b/pkg/sink/kafka/sarama_admin_test.go @@ -344,6 +344,7 @@ func TestIsUnretryableKafkaError(t *testing.T) { {name: "unsupported version", err: sarama.ErrUnsupportedVersion, unretryable: true}, {name: "invalid request", err: sarama.ErrInvalidRequest, unretryable: true}, {name: "client configuration error", err: sarama.ConfigurationError("invalid client config"), unretryable: true}, + {name: "wrapped client configuration error", err: errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ConfigurationError("invalid client config"), "describe-topic", "test-topic"), unretryable: true}, { name: "wrapped invalid topic", err: errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ErrInvalidTopic, "describe-topic", "test-topic"),