diff --git a/cdc/cdc/model/owner.go b/cdc/cdc/model/owner.go index 791d3b18..26e35136 100644 --- a/cdc/cdc/model/owner.go +++ b/cdc/cdc/model/owner.go @@ -165,6 +165,8 @@ type KeySpanOperation struct { // if the operation is a add operation, BoundaryTs is start ts BoundaryTs uint64 `json:"boundary_ts"` Status uint64 `json:"status,omitempty"` + + RelatedKeySpans []KeySpanLocation `json:"related_key_spans"` } // KeySpanProcessed returns whether the keyspan has been processed by processor @@ -217,7 +219,6 @@ type KeySpanReplicaInfo struct { StartTs Ts `json:"start-ts"` Start []byte End []byte - // MarkKeySpanID KeySpanID `json:"mark-keyspan-id"` } // Clone clones a KeySpanReplicaInfo @@ -270,7 +271,7 @@ func (ts *TaskStatus) RemoveKeySpan(id KeySpanID, boundaryTs Ts, isMoveKeySpan b } // AddKeySpan add the keyspan in KeySpanInfos and add a add kyespan operation. -func (ts *TaskStatus) AddKeySpan(id KeySpanID, keyspan *KeySpanReplicaInfo, boundaryTs Ts) { +func (ts *TaskStatus) AddKeySpan(id KeySpanID, keyspan *KeySpanReplicaInfo, boundaryTs Ts, relatedKeySpans []KeySpanLocation) { if ts.KeySpans == nil { ts.KeySpans = make(map[KeySpanID]*KeySpanReplicaInfo) } @@ -284,9 +285,10 @@ func (ts *TaskStatus) AddKeySpan(id KeySpanID, keyspan *KeySpanReplicaInfo, boun ts.Operation = make(map[KeySpanID]*KeySpanOperation) } ts.Operation[id] = &KeySpanOperation{ - Delete: false, - BoundaryTs: boundaryTs, - Status: OperDispatched, + Delete: false, + BoundaryTs: boundaryTs, + Status: OperDispatched, + RelatedKeySpans: relatedKeySpans, } } @@ -448,3 +450,9 @@ type ProcInfoSnap struct { CaptureID string `json:"capture-id"` KeySpans map[KeySpanID]*KeySpanReplicaInfo `json:"-"` } + +// KeySpanLocation records which capture a keyspan is in +type KeySpanLocation struct { + CaptureID string `json:"capture_id"` + KeySpanID KeySpanID `json:"keyspan_id"` +} diff --git a/cdc/cdc/model/owner_test.go b/cdc/cdc/model/owner_test.go index 2656eb89..77d303ac 100644 --- a/cdc/cdc/model/owner_test.go +++ b/cdc/cdc/model/owner_test.go @@ -48,21 +48,6 @@ func TestAdminJobType(t *testing.T) { } } -func TestDDLStateString(t *testing.T) { - t.Parallel() - - names := map[ChangeFeedDDLState]string{ - ChangeFeedSyncDML: "SyncDML", - ChangeFeedWaitToExecDDL: "WaitToExecDDL", - ChangeFeedExecDDL: "ExecDDL", - ChangeFeedDDLExecuteFailed: "DDLExecuteFailed", - ChangeFeedDDLState(100): "Unknown", - } - for state, name := range names { - require.Equal(t, name, state.String()) - } -} - func TestTaskPositionMarshal(t *testing.T) { t.Parallel() @@ -265,11 +250,11 @@ func TestAddKeySpan(t *testing.T) { }, } status := &TaskStatus{} - status.AddKeySpan(1, &KeySpanReplicaInfo{StartTs: ts}, ts) + status.AddKeySpan(1, &KeySpanReplicaInfo{StartTs: ts}, ts, nil) require.Equal(t, expected, status) // add existing keyspan does nothing - status.AddKeySpan(1, &KeySpanReplicaInfo{StartTs: 1}, 1) + status.AddKeySpan(1, &KeySpanReplicaInfo{StartTs: 1}, 1, nil) require.Equal(t, expected, status) } @@ -279,8 +264,8 @@ func TestTaskStatusApplyState(t *testing.T) { ts1 := uint64(420875042036766723) ts2 := uint64(420876783269969921) status := &TaskStatus{} - status.AddKeySpan(1, &KeySpanReplicaInfo{StartTs: ts1}, ts1) - status.AddKeySpan(2, &KeySpanReplicaInfo{StartTs: ts2}, ts2) + status.AddKeySpan(1, &KeySpanReplicaInfo{StartTs: ts1}, ts1, nil) + status.AddKeySpan(2, &KeySpanReplicaInfo{StartTs: ts2}, ts2, nil) require.True(t, status.SomeOperationsUnapplied()) require.Equal(t, ts1, status.AppliedTs()) diff --git a/cdc/cdc/owner/scheduler_v1.go b/cdc/cdc/owner/scheduler_v1.go index d492cd30..cde323ea 100644 --- a/cdc/cdc/owner/scheduler_v1.go +++ b/cdc/cdc/owner/scheduler_v1.go @@ -46,6 +46,8 @@ type schedulerJob struct { // if the operation is an add operation, boundaryTs is start ts BoundaryTs uint64 TargetCapture model.CaptureID + + RelatedKeySpans []model.KeySpanLocation } type moveKeySpanJob struct { @@ -55,7 +57,7 @@ type moveKeySpanJob struct { type oldScheduler struct { state *orchestrator.ChangefeedReactorState - currentKeySpansID []model.KeySpanID + currentKeySpanIDs []model.KeySpanID currentKeySpans map[model.KeySpanID]regionspan.Span captures map[model.CaptureID]*model.CaptureInfo @@ -80,20 +82,21 @@ func newSchedulerV1(f updateCurrentKeySpansFunc) scheduler { func (s *oldScheduler) Tick( ctx cdcContext.Context, state *orchestrator.ChangefeedReactorState, - // currentKeySpans []model.KeySpanID, captures map[model.CaptureID]*model.CaptureInfo, ) (shouldUpdateState bool, err error) { - s.state = state s.captures = captures - s.currentKeySpansID, s.currentKeySpans, err = s.updateCurrentKeySpans(ctx) + currentKeySpanIDs, currentKeySpans, err := s.updateCurrentKeySpans(ctx) if err != nil { return false, errors.Trace(err) } + newKeySpans, needRemoveKeySpans := s.diffCurrentKeySpans(currentKeySpans) + s.currentKeySpanIDs, s.currentKeySpans = currentKeySpanIDs, currentKeySpans + s.cleanUpFinishedOperations() - pendingJob, err := s.syncKeySpansWithCurrentKeySpans() + pendingJob, err := s.syncKeySpansWithCurrentKeySpans(newKeySpans, needRemoveKeySpans) if err != nil { return false, errors.Trace(err) } @@ -116,6 +119,27 @@ func (s *oldScheduler) Tick( return shouldUpdateState, nil } +func (s *oldScheduler) diffCurrentKeySpans(currentKeySpans map[model.KeySpanID]regionspan.Span) (map[model.KeySpanID]struct{}, []model.KeySpanID) { + oldKeySpans := s.currentKeySpans + + newKeySpans := map[model.KeySpanID]struct{}{} + needRemoveKeySpans := []model.KeySpanID{} + + for keyspanID := range oldKeySpans { + if _, ok := currentKeySpans[keyspanID]; !ok { + needRemoveKeySpans = append(needRemoveKeySpans, keyspanID) + } + } + + for keyspanID := range currentKeySpans { + if _, ok := oldKeySpans[keyspanID]; !ok { + newKeySpans[keyspanID] = struct{}{} + } + } + + return newKeySpans, needRemoveKeySpans +} + func (s *oldScheduler) MoveKeySpan(keyspanID model.KeySpanID, target model.CaptureID) { s.moveKeySpanJobQueue = append(s.moveKeySpanJobQueue, &moveKeySpanJob{ keyspanID: keyspanID, @@ -221,7 +245,9 @@ func (s *oldScheduler) dispatchToTargetCaptures(pendingJobs []*schedulerJob) { } } + count := 0 getMinWorkloadCapture := func() model.CaptureID { + count++ minCapture := "" minWorkLoad := uint64(math.MaxUint64) for captureID, workload := range workloads { @@ -249,26 +275,41 @@ func (s *oldScheduler) dispatchToTargetCaptures(pendingJobs []*schedulerJob) { // syncKeySpansWithCurrentKeySpans iterates all current keyspans to check whether it should be listened or not. // this function will return schedulerJob to make sure all keyspans will be listened. -func (s *oldScheduler) syncKeySpansWithCurrentKeySpans() ([]*schedulerJob, error) { +func (s *oldScheduler) syncKeySpansWithCurrentKeySpans(newKeySpans map[model.KeySpanID]struct{}, needRemoveKeySpans []model.KeySpanID) ([]*schedulerJob, error) { var pendingJob []*schedulerJob allKeySpanListeningNow, err := s.keyspan2CaptureIndex() if err != nil { return nil, errors.Trace(err) } + relatedKeySpans := make([]model.KeySpanLocation, 0, len(needRemoveKeySpans)) + for _, keyspanID := range needRemoveKeySpans { + if captureID, ok := allKeySpanListeningNow[keyspanID]; ok { + location := model.KeySpanLocation{ + CaptureID: captureID, + KeySpanID: keyspanID, + } + relatedKeySpans = append(relatedKeySpans, location) + } + } + globalCheckpointTs := s.state.Status.CheckpointTs - for _, keyspanID := range s.currentKeySpansID { + for _, keyspanID := range s.currentKeySpanIDs { if _, exist := allKeySpanListeningNow[keyspanID]; exist { delete(allKeySpanListeningNow, keyspanID) continue } - // For each keyspan which should be listened but is not, add an adding-keyspan job to the pending job list - pendingJob = append(pendingJob, &schedulerJob{ + job := &schedulerJob{ Tp: schedulerJobTypeAddKeySpan, KeySpanID: keyspanID, Start: s.currentKeySpans[keyspanID].Start, End: s.currentKeySpans[keyspanID].End, BoundaryTs: globalCheckpointTs, - }) + } + if _, ok := newKeySpans[keyspanID]; ok { + job.RelatedKeySpans = relatedKeySpans + } + // For each keyspan which should be listened but is not, add an adding-keyspan job to the pending job list + pendingJob = append(pendingJob, job) } // The remaining keyspans are the keyspans which should be not listened keyspansThatShouldNotBeListened := allKeySpanListeningNow @@ -303,7 +344,7 @@ func (s *oldScheduler) handleJobs(jobs []*schedulerJob) { StartTs: job.BoundaryTs, Start: job.Start, End: job.End, - }, job.BoundaryTs) + }, job.BoundaryTs, job.RelatedKeySpans) case schedulerJobTypeRemoveKeySpan: failpoint.Inject("OwnerRemoveKeySpanError", func() { // just skip removing this keyspan @@ -517,6 +558,7 @@ func (w *schedulerV1CompatWrapper) calculateWatermarks( if resolvedTs == model.Ts(math.MaxUint64) { return schedulerv2.CheckpointCannotProceed, 0 } + checkpointTs := resolvedTs for _, position := range state.TaskPositions { if checkpointTs > position.CheckPointTs { diff --git a/cdc/cdc/owner/scheduler_v1_test.go b/cdc/cdc/owner/scheduler_v1_test.go index e4741127..7fdbc651 100644 --- a/cdc/cdc/owner/scheduler_v1_test.go +++ b/cdc/cdc/owner/scheduler_v1_test.go @@ -16,6 +16,7 @@ package owner import ( "fmt" "math/rand" + "sort" "testing" "github.com/pingcap/check" @@ -141,21 +142,19 @@ func (s *schedulerSuite) TestScheduleOneCapture(c *check.C) { }) c.Assert(s.state.TaskStatuses[captureID].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ - 1: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, - 2: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, - 3: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, - 4: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, + 1: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, + 2: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, + 3: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, + 4: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, }) shouldUpdateState, err = s.scheduler.Tick(ctx, s.state, s.captures) // []model.KeySpanID{1, 2, 3, 4}, - c.Assert(err, check.IsNil) c.Assert(shouldUpdateState, check.IsTrue) s.tester.MustApplyPatches() // two keyspans finish adding operation s.finishKeySpanOperation(captureID, 2, 3) - s.scheduler.updateCurrentKeySpans = func(ctx cdcContext.Context) ([]model.KeySpanID, map[model.KeySpanID]regionspan.Span, error) { return []model.KeySpanID{3, 4, 5}, map[model.KeySpanID]regionspan.Span{ 3: {Start: []byte{'3'}, End: []byte{'4'}}, @@ -173,11 +172,20 @@ func (s *schedulerSuite) TestScheduleOneCapture(c *check.C) { 4: {StartTs: 0, Start: []byte{'4'}, End: []byte{'5'}}, 5: {StartTs: 0, Start: []byte{'5'}, End: []byte{'6'}}, }) + + keyspanOperation, IsTrue := s.state.TaskStatuses[captureID].Operation[5] + c.Assert(IsTrue, check.IsTrue) + sort.SliceStable(keyspanOperation.RelatedKeySpans, func(i, j int) bool { + return keyspanOperation.RelatedKeySpans[i].KeySpanID < keyspanOperation.RelatedKeySpans[j].KeySpanID + }) c.Assert(s.state.TaskStatuses[captureID].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ - 1: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched}, - 2: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched}, - 4: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, - 5: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, + 1: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 2: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 4: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, + 5: {Delete: false, + BoundaryTs: 0, + Status: model.OperDispatched, + RelatedKeySpans: []model.KeySpanLocation{{CaptureID: captureID, KeySpanID: 1}, {CaptureID: captureID, KeySpanID: 2}}}, }) // move a non exist keyspan to a non exist capture @@ -193,12 +201,21 @@ func (s *schedulerSuite) TestScheduleOneCapture(c *check.C) { 4: {StartTs: 0, Start: []byte{'4'}, End: []byte{'5'}}, 5: {StartTs: 0, Start: []byte{'5'}, End: []byte{'6'}}, }) + + keyspanOperation, IsTrue = s.state.TaskStatuses[captureID].Operation[5] + c.Assert(IsTrue, check.IsTrue) + sort.SliceStable(keyspanOperation.RelatedKeySpans, func(i, j int) bool { + return keyspanOperation.RelatedKeySpans[i].KeySpanID < keyspanOperation.RelatedKeySpans[j].KeySpanID + }) c.Assert(s.state.TaskStatuses[captureID].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ - 1: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched}, - 2: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched}, - 3: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched}, - 4: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, - 5: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, + 1: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 2: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 3: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 4: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, + 5: {Delete: false, + BoundaryTs: 0, + Status: model.OperDispatched, + RelatedKeySpans: []model.KeySpanLocation{{CaptureID: captureID, KeySpanID: 1}, {CaptureID: captureID, KeySpanID: 2}}}, }) // finish all operations @@ -236,7 +253,7 @@ func (s *schedulerSuite) TestScheduleOneCapture(c *check.C) { 5: {StartTs: 0, Start: []byte{'5'}, End: []byte{'6'}}, }) c.Assert(s.state.TaskStatuses[captureID].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ - 3: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, + 3: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, }) } @@ -265,7 +282,7 @@ func (s *schedulerSuite) TestScheduleMoveKeySpan(c *check.C) { 1: {StartTs: 0, Start: []byte{'1'}, End: []byte{'2'}}, }) c.Assert(s.state.TaskStatuses[captureID1].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ - 1: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, + 1: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, }) s.finishKeySpanOperation(captureID1, 1) @@ -295,7 +312,7 @@ func (s *schedulerSuite) TestScheduleMoveKeySpan(c *check.C) { 2: {StartTs: 0, Start: []byte{'2'}, End: []byte{'3'}}, }) c.Assert(s.state.TaskStatuses[captureID2].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ - 2: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched}, + 2: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, }) s.finishKeySpanOperation(captureID2, 2) @@ -311,7 +328,7 @@ func (s *schedulerSuite) TestScheduleMoveKeySpan(c *check.C) { c.Assert(s.state.TaskStatuses[captureID1].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{}) c.Assert(s.state.TaskStatuses[captureID2].KeySpans, check.DeepEquals, map[model.KeySpanID]*model.KeySpanReplicaInfo{}) c.Assert(s.state.TaskStatuses[captureID2].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ - 2: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched}, + 2: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, }) s.finishKeySpanOperation(captureID2, 2) @@ -429,3 +446,120 @@ func (s *schedulerSuite) TestScheduleRebalance(c *check.C) { } c.Assert(keyspanIDs, check.DeepEquals, map[model.KeySpanID]struct{}{1: {}, 2: {}, 3: {}, 4: {}, 5: {}, 6: {}}) } + +func (s *schedulerSuite) TestRelatedKeySpans(c *check.C) { + defer testleak.AfterTest(c)() + s.reset(c) + captureID := "test-capture" + s.addCapture(captureID) + + ctx := cdcContext.NewBackendContext4Test(false) + ctx, cancel := cdcContext.WithCancel(ctx) + defer cancel() + + s.scheduler.updateCurrentKeySpans = func(ctx cdcContext.Context) ([]model.KeySpanID, map[model.KeySpanID]regionspan.Span, error) { + return []model.KeySpanID{1}, map[model.KeySpanID]regionspan.Span{ + 1: {Start: []byte{'1'}, End: []byte{'3'}}, + }, nil + } + + shouldUpdateState, err := s.scheduler.Tick(ctx, s.state, s.captures) // []model.KeySpanID{1}, + c.Assert(err, check.IsNil) + c.Assert(shouldUpdateState, check.IsFalse) + s.tester.MustApplyPatches() + c.Assert(s.state.TaskStatuses[captureID].KeySpans, check.DeepEquals, map[model.KeySpanID]*model.KeySpanReplicaInfo{ + 1: {StartTs: 0, Start: []byte{'1'}, End: []byte{'3'}}, + }) + c.Assert(s.state.TaskStatuses[captureID].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ + 1: {Delete: false, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{}}, + }) + + s.state.PatchTaskWorkload(captureID, func(workload model.TaskWorkload) (model.TaskWorkload, bool, error) { + if workload == nil { + workload = make(model.TaskWorkload) + } + for keyspanID := range s.state.TaskStatuses[captureID].KeySpans { + if s.state.TaskStatuses[captureID].Operation[keyspanID].Delete { + delete(workload, keyspanID) + } else { + workload[keyspanID] = model.WorkloadInfo{ + Workload: 1, + } + } + } + return workload, true, nil + }) + s.tester.MustApplyPatches() + + s.scheduler.updateCurrentKeySpans = func(ctx cdcContext.Context) ([]model.KeySpanID, map[model.KeySpanID]regionspan.Span, error) { + return []model.KeySpanID{2, 3}, map[model.KeySpanID]regionspan.Span{ + 2: {Start: []byte{'1'}, End: []byte{'2'}}, 3: {Start: []byte{'2'}, End: []byte{'3'}}, + }, nil + } + + shouldUpdateState, err = s.scheduler.Tick(ctx, s.state, s.captures) // []model.KeySpanID{2, 3}, + c.Assert(err, check.IsNil) + c.Assert(shouldUpdateState, check.IsFalse) + s.tester.MustApplyPatches() + c.Assert(s.state.TaskStatuses[captureID].KeySpans, check.DeepEquals, map[model.KeySpanID]*model.KeySpanReplicaInfo{ + 2: {StartTs: 0, Start: []byte{'1'}, End: []byte{'2'}}, + 3: {StartTs: 0, Start: []byte{'2'}, End: []byte{'3'}}, + }) + c.Assert(s.state.TaskStatuses[captureID].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ + 1: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 2: {Delete: false, + BoundaryTs: 0, + Status: model.OperDispatched, + RelatedKeySpans: []model.KeySpanLocation{{CaptureID: captureID, KeySpanID: 1}}}, + 3: {Delete: false, + BoundaryTs: 0, + Status: model.OperDispatched, + RelatedKeySpans: []model.KeySpanLocation{{CaptureID: captureID, KeySpanID: 1}}}, + }) + + s.state.PatchTaskWorkload(captureID, func(workload model.TaskWorkload) (model.TaskWorkload, bool, error) { + if workload == nil { + workload = make(model.TaskWorkload) + } + for keyspanID := range s.state.TaskStatuses[captureID].KeySpans { + if s.state.TaskStatuses[captureID].Operation[keyspanID].Delete { + delete(workload, keyspanID) + } else { + workload[keyspanID] = model.WorkloadInfo{ + Workload: 1, + } + } + } + return workload, true, nil + }) + s.tester.MustApplyPatches() + + s.scheduler.updateCurrentKeySpans = func(ctx cdcContext.Context) ([]model.KeySpanID, map[model.KeySpanID]regionspan.Span, error) { + return []model.KeySpanID{4}, map[model.KeySpanID]regionspan.Span{ + 4: {Start: []byte{'1'}, End: []byte{'3'}}, + }, nil + } + shouldUpdateState, err = s.scheduler.Tick(ctx, s.state, s.captures) // []model.KeySpanID{4}, + c.Assert(err, check.IsNil) + c.Assert(shouldUpdateState, check.IsFalse) + s.tester.MustApplyPatches() + c.Assert(s.state.TaskStatuses[captureID].KeySpans, check.DeepEquals, map[model.KeySpanID]*model.KeySpanReplicaInfo{ + 4: {StartTs: 0, Start: []byte{'1'}, End: []byte{'3'}}, + }) + + keyspanOperation, IsTrue := s.state.TaskStatuses[captureID].Operation[4] + c.Assert(IsTrue, check.IsTrue) + sort.SliceStable(keyspanOperation.RelatedKeySpans, func(i, j int) bool { + return keyspanOperation.RelatedKeySpans[i].KeySpanID < keyspanOperation.RelatedKeySpans[j].KeySpanID + }) + + c.Assert(s.state.TaskStatuses[captureID].Operation, check.DeepEquals, map[model.KeySpanID]*model.KeySpanOperation{ + 1: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 2: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 3: {Delete: true, BoundaryTs: 0, Status: model.OperDispatched, RelatedKeySpans: nil}, + 4: {Delete: false, + BoundaryTs: 0, + Status: model.OperDispatched, + RelatedKeySpans: []model.KeySpanLocation{{CaptureID: captureID, KeySpanID: 2}, {CaptureID: captureID, KeySpanID: 3}}}, + }) +} diff --git a/cdc/cdc/processor/processor.go b/cdc/cdc/processor/processor.go index 6bb1b679..f0bdc920 100644 --- a/cdc/cdc/processor/processor.go +++ b/cdc/cdc/processor/processor.go @@ -558,6 +558,10 @@ func (p *processor) handleKeySpanOperation(ctx cdcContext.Context) error { if replicaInfo.StartTs != opt.BoundaryTs { log.Warn("the startTs and BoundaryTs of add keyspan operation should be always equaled", zap.Any("replicaInfo", replicaInfo)) } + + if !p.checkRelatedKeyspans(opt.RelatedKeySpans) { + continue + } err := p.addKeySpan(ctx, keyspanID, replicaInfo) if err != nil { return errors.Trace(err) @@ -597,6 +601,18 @@ func (p *processor) handleKeySpanOperation(ctx cdcContext.Context) error { return nil } +func (p *processor) checkRelatedKeyspans(relatedKeySpans []model.KeySpanLocation) bool { + for _, location := range relatedKeySpans { + if taskStatus, ok := p.changefeed.TaskStatuses[location.CaptureID]; ok { + if operation, ok := taskStatus.Operation[location.KeySpanID]; ok && operation.Status != model.OperFinished { + return false + } + } + + } + return true +} + func (p *processor) sendError(err error) { if err == nil { return diff --git a/cdc/cdc/processor/processor_test.go b/cdc/cdc/processor/processor_test.go index c225d79a..faeab9c2 100644 --- a/cdc/cdc/processor/processor_test.go +++ b/cdc/cdc/processor/processor_test.go @@ -217,7 +217,7 @@ func (s *processorSuite) TestHandleKeySpanOperation4SingleKeySpan(c *check.C) { // add keyspan, in processing // in current implementation of owner, the startTs and BoundaryTs of add keyspan operation should be always equaled. p.changefeed.PatchTaskStatus(p.captureInfo.ID, func(status *model.TaskStatus) (*model.TaskStatus, bool, error) { - status.AddKeySpan(66, &model.KeySpanReplicaInfo{StartTs: 60}, 60) + status.AddKeySpan(66, &model.KeySpanReplicaInfo{StartTs: 60}, 60, nil) return status, true, nil }) tester.MustApplyPatches() @@ -351,9 +351,9 @@ func (s *processorSuite) TestHandleKeySpanOperation4MultiKeySpan(c *check.C) { // add keyspan, in processing // in current implementation of owner, the startTs and BoundaryTs of add keyspan operation should be always equaled. p.changefeed.PatchTaskStatus(p.captureInfo.ID, func(status *model.TaskStatus) (*model.TaskStatus, bool, error) { - status.AddKeySpan(1, &model.KeySpanReplicaInfo{StartTs: 60}, 60) - status.AddKeySpan(2, &model.KeySpanReplicaInfo{StartTs: 50}, 50) - status.AddKeySpan(3, &model.KeySpanReplicaInfo{StartTs: 40}, 40) + status.AddKeySpan(1, &model.KeySpanReplicaInfo{StartTs: 60}, 60, nil) + status.AddKeySpan(2, &model.KeySpanReplicaInfo{StartTs: 50}, 50, nil) + status.AddKeySpan(3, &model.KeySpanReplicaInfo{StartTs: 40}, 40, nil) status.KeySpans[4] = &model.KeySpanReplicaInfo{StartTs: 30} return status, true, nil }) @@ -926,3 +926,80 @@ func (s *processorSuite) TestIgnorableError(c *check.C) { c.Assert(isProcessorIgnorableError(tc.err), check.Equals, tc.ignorable) } } + +func (s *processorSuite) TestHandleKeySpanOperationWithRelatedKeySpans(c *check.C) { + defer testleak.AfterTest(c)() + ctx := cdcContext.NewBackendContext4Test(true) + p, tester := initProcessor4Test(ctx, c) + var err error + + // no operation + _, err = p.Tick(ctx, p.changefeed) + c.Assert(err, check.IsNil) + tester.MustApplyPatches() + + // add keyspan1, in processing + p.changefeed.PatchTaskStatus(p.captureInfo.ID, func(status *model.TaskStatus) (*model.TaskStatus, bool, error) { + status.AddKeySpan(1, &model.KeySpanReplicaInfo{StartTs: 60}, 80, nil) + return status, true, nil + }) + tester.MustApplyPatches() + _, err = p.Tick(ctx, p.changefeed) + c.Assert(err, check.IsNil) + tester.MustApplyPatches() + c.Assert(p.changefeed.TaskStatuses[p.captureInfo.ID], check.DeepEquals, &model.TaskStatus{ + KeySpans: map[uint64]*model.KeySpanReplicaInfo{ + 1: {StartTs: 60}, + }, + Operation: map[uint64]*model.KeySpanOperation{ + 1: {Delete: false, BoundaryTs: 80, Status: model.OperProcessed}, + }, + }) + c.Assert(p.keyspans, check.HasLen, 1) + c.Assert(p.changefeed.TaskPositions[p.captureInfo.ID].CheckPointTs, check.Equals, uint64(60)) + c.Assert(p.changefeed.TaskPositions[p.captureInfo.ID].ResolvedTs, check.Equals, uint64(60)) + + // add keyspan2 & keyspan3, remove keyspan1 + p.changefeed.PatchTaskStatus(p.captureInfo.ID, func(status *model.TaskStatus) (*model.TaskStatus, bool, error) { + status.AddKeySpan(2, &model.KeySpanReplicaInfo{StartTs: 60}, 60, []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}) + status.AddKeySpan(3, &model.KeySpanReplicaInfo{StartTs: 60}, 60, []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}) + status.RemoveKeySpan(1, 60, false) + return status, true, nil + }) + tester.MustApplyPatches() + // try to stop keyspand1 + _, err = p.Tick(ctx, p.changefeed) + c.Assert(err, check.IsNil) + tester.MustApplyPatches() + c.Assert(p.changefeed.TaskStatuses[p.captureInfo.ID].KeySpans, check.DeepEquals, map[uint64]*model.KeySpanReplicaInfo{ + 2: {StartTs: 60}, + 3: {StartTs: 60}, + }) + c.Assert(p.changefeed.TaskStatuses[p.captureInfo.ID].Operation, check.DeepEquals, map[uint64]*model.KeySpanOperation{ + 1: {Delete: true, BoundaryTs: 60, Status: model.OperProcessed}, + 2: {Delete: false, BoundaryTs: 60, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}}, + 3: {Delete: false, BoundaryTs: 60, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}}, + }) + keyspan1 := p.keyspans[1].(*mockKeySpanPipeline) + keyspan1.status = keyspanpipeline.KeySpanStatusStopped + + // finish stoping keyspand1 + _, err = p.Tick(ctx, p.changefeed) + c.Assert(err, check.IsNil) + tester.MustApplyPatches() + c.Assert(p.changefeed.TaskStatuses[p.captureInfo.ID].Operation, check.DeepEquals, map[uint64]*model.KeySpanOperation{ + 1: {Delete: true, BoundaryTs: 60, Status: model.OperFinished}, + 2: {Delete: false, BoundaryTs: 60, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}}, + 3: {Delete: false, BoundaryTs: 60, Status: model.OperDispatched, RelatedKeySpans: []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}}, + }) + cleanUpFinishedOpOperation(p.changefeed, p.captureInfo.ID, tester) + + // start keyspan2 & keyspan3 + _, err = p.Tick(ctx, p.changefeed) + c.Assert(err, check.IsNil) + tester.MustApplyPatches() + c.Assert(p.changefeed.TaskStatuses[p.captureInfo.ID].Operation, check.DeepEquals, map[uint64]*model.KeySpanOperation{ + 2: {Delete: false, BoundaryTs: 60, Status: model.OperProcessed, RelatedKeySpans: []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}}, + 3: {Delete: false, BoundaryTs: 60, Status: model.OperProcessed, RelatedKeySpans: []model.KeySpanLocation{{CaptureID: p.captureInfo.ID, KeySpanID: 1}}}, + }) +}