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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 43 additions & 0 deletions maintainer/maintainer_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -259,6 +259,49 @@
continue
}
spanController.UpdateStatus(stm, status)
<<<<<<< HEAD

Check failure on line 262 in maintainer/maintainer_controller.go

View workflow job for this annotation

GitHub Actions / Build Classic CDC

syntax error: unexpected <<, expected }
=======

Check failure on line 263 in maintainer/maintainer_controller.go

View workflow job for this annotation

GitHub Actions / Build Classic CDC

syntax error: unexpected ==, expected }

if !allowSelfHealing {
continue
}

// Fallback: dispatcher becomes non-working without an operator.
//
// In normal scheduling flow, a dispatcher should transition to Stopped/Removed as part of a maintainer
// operator (Remove/Move/Split...). However, after maintainer failover we can lose operatorController state
// while dispatcher managers keep executing the already-issued requests.
//
// A real example is a "remove request in transit" during bootstrap:
// - Old maintainer sends a Remove (e.g. the remove-origin phase of Move), but the request hasn't reached
// dispatcher manager yet.
// - New maintainer bootstraps from dispatcher manager snapshots and sees the dispatcher as Working, with
// no in-flight operator reported in bootstrap response.
// - After bootstrap, the in-transit Remove arrives, the dispatcher is removed, and the new maintainer
// observes a terminal status without a corresponding operator.
//
// In these cases we'd observe a non-working status but have no operator to drive the follow-up
// rescheduling, so we mark the span absent to let the scheduler recreate it.
//
// Safety against message reordering/resend:
// MarkSpanAbsentIfCurrent atomically verifies that stm is still the current desired task and is still bound
// to the reporting node. A concurrent split, merge, move, or DDL removal therefore makes this a no-op.
if status.ComponentStatus == heartbeatpb.ComponentState_Stopped ||
status.ComponentStatus == heartbeatpb.ComponentState_Removed {
if op := operatorController.GetOperator(dispatcherID); op == nil {
if c.removeTerminalSpanCoveredByMergedSpan(spanController, stm) {
continue
}
if spanController.MarkSpanAbsentIfCurrent(stm, from) {
log.Warn("dispatcher becomes non-working without operator, mark span absent for rescheduling",
zap.String("changefeed", c.changefeedID.Name()),
zap.String("from", from.String()),
zap.String("dispatcherID", dispatcherID.String()),
zap.Any("status", status))
}
}
}
>>>>>>> 3f0a68abf (maintainer: prevent stale spans from reentering scheduler state (#6073))

Check failure on line 304 in maintainer/maintainer_controller.go

View workflow job for this annotation

GitHub Actions / Build Classic CDC

invalid character U+0023 '#'
}
}

Expand Down
3 changes: 3 additions & 0 deletions maintainer/operator/operator_move_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ func setupTestEnvironment(t *testing.T) (*span.Controller, common.ChangeFeedID,
// 4. Verify that the move is aborted and the span is marked absent after origin is stopped
func TestMoveOperator_DestNodeRemovedBeforeOriginStopped(t *testing.T) {
spanController, _, replicaSet, nodeA, nodeB := setupTestEnvironment(t)
spanController.AddReplicatingSpan(replicaSet)

op := NewMoveDispatcherOperator(spanController, replicaSet, nodeA, nodeB, 7)
require.NotNil(t, op)
Expand Down Expand Up @@ -131,6 +132,7 @@ func TestMoveOperator_DestNodeRemovedBeforeOriginStopped(t *testing.T) {
// 5. Verify that the span is marked as absent for rescheduling
func TestMoveOperator_DestNodeRemovedAfterOriginStopped(t *testing.T) {
spanController, _, replicaSet, nodeA, nodeB := setupTestEnvironment(t)
spanController.AddReplicatingSpan(replicaSet)

op := NewMoveDispatcherOperator(spanController, replicaSet, nodeA, nodeB, 7)
require.NotNil(t, op)
Expand Down Expand Up @@ -270,6 +272,7 @@ func TestMoveOperator_BothNodesRemovedBeforeStartDoesNotLeaveSchedulingWithoutNo
// 4. Verify that the move is aborted and the span becomes absent for rescheduling
func TestMoveOperator_DestThenOriginRemovedAbortsToAbsent(t *testing.T) {
spanController, _, replicaSet, nodeA, nodeB := setupTestEnvironment(t)
spanController.AddReplicatingSpan(replicaSet)

op := NewMoveDispatcherOperator(spanController, replicaSet, nodeA, nodeB, 7)
require.NotNil(t, op)
Expand Down
2 changes: 2 additions & 0 deletions maintainer/operator/operator_split_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ import (
// 4. Verify that the span is marked as absent, and split is not executed
func TestSplitOperator_OriginNodeRemovedBeforeStopped(t *testing.T) {
spanController, _, replicaSet, nodeA, _ := setupTestEnvironment(t)
spanController.AddReplicatingSpan(replicaSet)

// Define split spans
splitSpans := []*heartbeatpb.TableSpan{
Expand Down Expand Up @@ -82,6 +83,7 @@ func TestSplitOperator_OriginNodeRemovedBeforeStopped(t *testing.T) {
// 4. Verify that the span is still marked as absent, and split is not executed
func TestSplitOperator_OriginNodeRemovedAfterStopped(t *testing.T) {
spanController, _, replicaSet, nodeA, _ := setupTestEnvironment(t)
spanController.AddReplicatingSpan(replicaSet)

// Define split spans
splitSpans := []*heartbeatpb.TableSpan{
Expand Down
37 changes: 33 additions & 4 deletions maintainer/span/span_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -394,12 +394,31 @@ func (c *Controller) AddReplicatingSpan(span *replica.SpanReplication) {
c.untrackNonReplicatingSpan(span)
}

// MarkSpanAbsent marks span as absent
func (c *Controller) MarkSpanAbsent(span *replica.SpanReplication) {
// MarkSpanAbsent marks span as absent if it is still the current task for its dispatcher ID.
func (c *Controller) MarkSpanAbsent(span *replica.SpanReplication) bool {
if span == nil {
return false
}

c.mu.Lock()
defer c.mu.Unlock()
c.MarkAbsentWithoutLock(span)
c.trackNonReplicatingSpan(span)
return c.markSpanAbsentIfCurrentWithoutLock(span)
}

// MarkSpanAbsentIfCurrent marks span as absent if it is still the current task
// and is still bound to expectedNode.
func (c *Controller) MarkSpanAbsentIfCurrent(span *replica.SpanReplication, expectedNode node.ID) bool {
if span == nil {
return false
}

c.mu.Lock()
defer c.mu.Unlock()
current, ok := c.allTasks[span.ID]
if !ok || current != span || current.GetNodeID() != expectedNode {
return false
}
return c.markSpanAbsentIfCurrentWithoutLock(span)
}

// MarkSpanScheduling marks span as scheduling
Expand Down Expand Up @@ -607,6 +626,16 @@ func (c *Controller) removeSpanWithoutLock(spans ...*replica.SpanReplication) {
}
}

func (c *Controller) markSpanAbsentIfCurrentWithoutLock(span *replica.SpanReplication) bool {
current, ok := c.allTasks[span.ID]
if !ok || current != span {
return false
}
c.MarkAbsentWithoutLock(current)
c.trackNonReplicatingSpan(current)
return true
}

func (c *Controller) trackNonReplicatingSpan(span *replica.SpanReplication) {
if span == c.ddlSpan {
return
Expand Down
52 changes: 51 additions & 1 deletion maintainer/span/span_controller_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -592,6 +592,32 @@ func TestReplaceReplicaSet(t *testing.T) {
require.Equal(t, 2, controller.GetTaskSizeBySchemaID(1))
}

func TestMarkSpanAbsentIgnoresRemovedSpan(t *testing.T) {
controller := newControllerWithCheckerForTest(t)
replicaSpanID := common.NewDispatcherID()
replicaSpan := replica.NewWorkingSpanReplication(controller.changefeedID, replicaSpanID,
1,
testutil.GetTableSpanByID(3), &heartbeatpb.TableSpanStatus{
ID: replicaSpanID.ToPB(),
ComponentStatus: heartbeatpb.ComponentState_Working,
CheckpointTs: 1,
}, "node1", false)
controller.AddReplicatingSpan(replicaSpan)

controller.ReplaceReplicaSet(
[]*replica.SpanReplication{replicaSpan},
[]*heartbeatpb.TableSpan{testutil.GetTableSpanByID(3), testutil.GetTableSpanByID(4)},
5,
[]node.ID{},
)
require.False(t, controller.MarkSpanAbsent(replicaSpan))

require.Nil(t, controller.GetTaskByID(replicaSpan.ID))
require.NotContains(t, controller.GetAbsentForTest(3), replicaSpan)
require.NotContains(t, controller.nonReplicatingCheckpointTs.checkpointTsBySpanID, replicaSpan.ID)
require.Equal(t, 2, controller.GetAbsentSize())
}

// TestMarkSpanAbsent tests the MarkSpanAbsent functionality
func TestMarkSpanAbsent(t *testing.T) {
controller := newControllerWithCheckerForTest(t)
Expand All @@ -605,8 +631,32 @@ func TestMarkSpanAbsent(t *testing.T) {
CheckpointTs: 1,
}, "node1", false)
controller.AddReplicatingSpan(replicaSpan)
controller.MarkSpanAbsent(replicaSpan)
require.True(t, controller.MarkSpanAbsent(replicaSpan))
require.Equal(t, 1, controller.GetAbsentSize())
require.Equal(t, "", replicaSpan.GetNodeID().String())
}

func TestMarkSpanAbsentIfCurrentRejectsOldOwner(t *testing.T) {
controller := newControllerWithCheckerForTest(t)
replicaSpanID := common.NewDispatcherID()
replicaSpan := replica.NewWorkingSpanReplication(controller.changefeedID, replicaSpanID,
1,
testutil.GetTableSpanByID(3), &heartbeatpb.TableSpanStatus{
ID: replicaSpanID.ToPB(),
ComponentStatus: heartbeatpb.ComponentState_Working,
CheckpointTs: 1,
}, "node1", false)
controller.AddReplicatingSpan(replicaSpan)
controller.BindSpanToNode("node1", "node2", replicaSpan)

require.False(t, controller.MarkSpanAbsentIfCurrent(replicaSpan, "node1"))
require.Equal(t, 0, controller.GetAbsentSize())
require.Equal(t, 1, controller.GetSchedulingSize())
require.Equal(t, "node2", replicaSpan.GetNodeID().String())

require.True(t, controller.MarkSpanAbsentIfCurrent(replicaSpan, "node2"))
require.Equal(t, 1, controller.GetAbsentSize())
require.Equal(t, 0, controller.GetSchedulingSize())
require.Equal(t, "", replicaSpan.GetNodeID().String())
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -181,7 +181,7 @@ main_with_consistent() {

kill -9 $NORMAL_TABLE_DDL_PID ${pids[@]} $MERGE_AND_SPLIT_TABLE_PID
# to ensure row changed events have been replicated to TiCDC
sleep 20
sleep 30
if ((RANDOM % 2)); then
# For rename table, modify column ddl, drop column, drop index and drop table ddl, the struct of table is wrong when appling snapshot.
# see https://github.com/pingcap/tidb/issues/63464.
Expand Down
4 changes: 4 additions & 0 deletions tests/integration_tests/run_heavy_it_in_ci.sh
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,10 @@ echo "Group Number (parsed): ${group_num}"
if [[ $group_num =~ ^[0-9]+$ ]] && [[ -n ${groups[10#${group_num}]} ]]; then
# force use decimal index
test_names="${groups[10#${group_num}]}"
if [[ "$sink_type" == "mysql" ]]; then
# Temporarily run the regression case in every MySQL shard.
test_names="ddl_for_split_tables_with_random_merge_and_split"
fi
# Run test cases
echo "Run cases: ${test_names}"
export TICDC_NEWARCH=true
Expand Down
Loading