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 caf9f324cb..c960f0cca1 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) } @@ -260,6 +263,7 @@ func (s *sink) WriteBlockEvent(event commonEvent.BlockEvent) error { } func (s *sink) close() { + s.isNormal.Store(false) s.eventChan.Close() s.rowChan.Close() } @@ -315,6 +319,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...) } } @@ -567,6 +576,7 @@ func (s *sink) getAllTableNames(ts uint64) []*commonEvent.SchemaTableName { } func (s *sink) Close() { + 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..b4de4edcbb 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -110,7 +110,7 @@ func TestVerifyInvalidConfig(t *testing.T) { factory := kafka.NewMockFactory(ctrl) gomock.InOrder( factory.EXPECT().AdminClient(gomock.Any()).Return(adminClient, nil), - adminClient.EXPECT().GetTopicsMeta([]string{kafkaSinkTestTopic}, true).Return( + adminClient.EXPECT().GetTopicsMeta([]string{kafkaSinkTestTopic}, false).Return( map[string]kafka.TopicDetail{kafkaSinkTestTopic: {Name: kafkaSinkTestTopic}}, nil), adminClient.EXPECT().Close(), ) @@ -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( + asyncProducer.EXPECT().Close(), + topicManager.EXPECT().Close(), + adminClient.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( + syncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }), + asyncProducer.EXPECT().Close().Do(func() { closeCount.Add(1) }), + topicManager.EXPECT().Close().Do(func() { closeCount.Add(1) }), + adminClient.EXPECT().Close().Do(func() { closeCount.Add(1) }), + ) kafkaSink, err := newWithComponents( t.Context(), @@ -313,8 +319,20 @@ func TestKafkaSinkConstructionAndCleanup(t *testing.T) { require.NoError(t, err) require.Zero(t, closeCount.Load()) + require.True(t, kafkaSink.IsNormal()) + 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) + require.False(t, ok) + _, ok, err = kafkaSink.rowChan.GetWithContext(t.Context()) + require.NoError(t, err) + require.False(t, ok) }) } @@ -394,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 a0110217f3..b6d6c277c7 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,9 @@ func GetTopicManagerAndTryCreateTopic( } ctx, cancel := context.WithCancel(ctx) topicManager.cancel = cancel - go topicManager.backgroundRefreshMeta(ctx) + topicManager.wg.Go(func() { + topicManager.backgroundRefreshMeta(ctx) + }) return topicManager, nil } @@ -214,6 +217,9 @@ func (m *kafkaTopicManager) waitUntilTopicVisible( }, retry.WithBackoffBaseDelay(500), retry.WithBackoffMaxDelay(1000), retry.WithMaxTries(6), + retry.WithIsRetryableErr(func(err error) bool { + return !kafka.IsUnretryableKafkaError(err) + }), ) if err != nil { log.Warn("kafka topic metadata refresh failed", @@ -260,8 +266,6 @@ func (m *kafkaTopicManager) createTopic( return 0, err } - m.tryUpdatePartitionsAndLogging(topicName, m.cfg.PartitionNum) - return m.cfg.PartitionNum, nil } @@ -272,27 +276,15 @@ func (m *kafkaTopicManager) createTopic( func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( ctx context.Context, topicName string, ) (int32, error) { - // If the topic is not in the cache, we try to get the metadata of the topic. - // ignoreTopicErr is set to true to ignore the error if the topic is not found, - // which means we should create the topic later. - topicDetails, err := m.admin.GetTopicsMeta([]string{topicName}, true) - if err != nil { - if kafka.IsAuthorizationFailed(err) { - return m.useConfiguredPartitionNum(topicName, err), nil + // If the topic is not in the cache, try to get its metadata. + topicDetails, err := m.admin.GetTopicsMeta([]string{topicName}, false) + if err == nil { + if numPartition, ok := m.tryStoreTopicMeta(topicName, topicDetails); ok { + return numPartition, nil } - return 0, err } - if numPartition, ok := m.tryStoreTopicMeta(topicName, topicDetails); ok { - return numPartition, nil - } - - topicDetails, err = m.admin.GetTopicsMeta([]string{topicName}, false) - if err != nil { - if kafka.IsAuthorizationFailed(err) { - return m.useConfiguredPartitionNum(topicName, err), nil - } - } else if numPartition, ok := m.tryStoreTopicMeta(topicName, topicDetails); ok { - return numPartition, nil + if kafka.IsAuthorizationFailed(err) { + return m.useConfiguredPartitionNum(topicName, err), nil } start := time.Now() @@ -308,6 +300,7 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( if err != nil { return 0, err } + m.tryUpdatePartitionsAndLogging(topicName, partitionNum) log.Info( "kafka topic created", @@ -348,7 +341,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 e4dc2612eb..09c44fceeb 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -15,8 +15,11 @@ package topicmanager import ( "context" + "io" "testing" + "time" + "github.com/IBM/sarama" "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/errors" @@ -36,12 +39,9 @@ func TestCreateTopic(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{kafkaTopicManagerTestTopic}, true).Return( + adminClient.EXPECT().GetTopicsMeta([]string{kafkaTopicManagerTestTopic}, false).Return( map[string]kafka.TopicDetail{ - kafkaTopicManagerTestTopic: { - Name: kafkaTopicManagerTestTopic, - NumPartitions: 2, - }, + kafkaTopicManagerTestTopic: {Name: kafkaTopicManagerTestTopic, NumPartitions: 2}, }, nil) manager := newKafkaTopicManager( kafkaTopicManagerTestTopic, @@ -62,26 +62,30 @@ func TestCreateTopic(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) var createdTopic *kafka.TopicDetail - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + postCreateDescribeCount := 0 + var manager *kafkaTopicManager 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: { - Name: createdTopic.Name, - NumPartitions: createdTopic.NumPartitions, - }, + createdTopic.Name: {Name: createdTopic.Name, NumPartitions: 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, @@ -102,6 +106,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) @@ -112,7 +117,6 @@ func TestCreateTopic(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) manager := newKafkaTopicManager( "new-topic", @@ -136,7 +140,6 @@ func TestCreateTopic(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) var createdTopic *kafka.TopicDetail adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( @@ -168,7 +171,6 @@ func TestCreateTopicValidatesReplicationFactor(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetBrokerConfig(kafka.MinInsyncReplicasConfigName).Return("2", true, nil) manager := newKafkaTopicManager( @@ -188,55 +190,26 @@ 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) { +func TestWaitUntilTopicVisibleUnretryableError(t *testing.T) { t.Parallel() ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - created := false - postCreateDescribeCount := 0 - 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++ - if postCreateDescribeCount == 1 { - return map[string]kafka.TopicDetail{}, nil - } - return map[string]kafka.TopicDetail{ - "delayed-topic": { - Name: "delayed-topic", - NumPartitions: 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 - }) - - err := EnsureTopic( - context.Background(), + 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"), - "delayed-topic", - &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, - ReplicationFactor: 1, - }, adminClient, + &kafka.AutoCreateTopicConfig{PartitionNum: 2}, ) - require.NoError(t, err) - require.Equal(t, 2, postCreateDescribeCount) + err := manager.waitUntilTopicVisible(context.Background(), "invalid-topic") + + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, sarama.ErrInvalidTopic) } func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { @@ -244,12 +217,9 @@ func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"existing-topic"}, true).Return( + adminClient.EXPECT().GetTopicsMeta([]string{"existing-topic"}, false).Return( map[string]kafka.TopicDetail{ - "existing-topic": { - Name: "existing-topic", - NumPartitions: 2, - }, + "existing-topic": {Name: "existing-topic", NumPartitions: 2}, }, nil) manager, err := GetTopicManagerAndTryCreateTopic( @@ -265,12 +235,47 @@ 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.Go(func() { + <-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() ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, false).Return( nil, errors.ErrKafkaAuthorizationFailed.GenWithStackByArgs("describe-topic", "default-topic")) manager := newKafkaTopicManager( @@ -298,7 +303,6 @@ func TestCreateTopicWithCreateDenied(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) adminClient.EXPECT().CreateTopic(&kafka.TopicDetail{ Name: "default-topic", diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index d45ec702d1..8c1a54f2ec 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) @@ -159,6 +156,26 @@ func IsAuthorizationFailed(err error) bool { errors.Is(err, sarama.ErrClusterAuthorizationFailed) } +// 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) || + 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) { result := make(map[string]int32, len(topics)) for _, topic := range topics { diff --git a/pkg/sink/kafka/options.go b/pkg/sink/kafka/options.go index 81cc997d6f..7f62babdc2 100644 --- a/pkg/sink/kafka/options.go +++ b/pkg/sink/kafka/options.go @@ -633,6 +633,8 @@ func adjustOptions( options *options, topic string, ) error { + // The topic may not exist yet and will be created later by the topic manager, + // so ignore per-topic metadata errors here. topics, err := admin.GetTopicsMeta([]string{topic}, true) if err != nil { return err diff --git a/pkg/sink/kafka/sarama_admin_test.go b/pkg/sink/kafka/sarama_admin_test.go index 90fc3530dd..d0ce07fc1d 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 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{ @@ -167,6 +167,31 @@ func TestGetTopicsMeta(t *testing.T) { topics, err := client.GetTopicsMeta([]string{"valid-topic", "missing-topic"}, false) + require.Nil(t, topics) + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + 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": { @@ -248,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{ @@ -287,6 +312,53 @@ func TestIsAuthorizationFailed(t *testing.T) { } } +func TestIsUnretryableKafkaError(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + err error + unretryable bool + }{ + {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 unknown 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}, + {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 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"), + unretryable: true, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + require.Equal(t, test.unretryable, IsUnretryableKafkaError(test.err)) + }) + } +} + func TestCreateTopic(t *testing.T) { t.Parallel()