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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 14 additions & 12 deletions pkg/cdc/sql_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -248,19 +248,20 @@ const (
"table_name = '%s'"

CDCCollectTableInfoSqlTemplate = "SELECT " +
" rel_id, " +
" relname, " +
" reldatabase_id, " +
" reldatabase, " +
" rel_createsql, " +
" account_id " +
"FROM `mo_catalog`.`mo_tables` " +
" tbl.rel_id, " +
" tbl.relname, " +
" tbl.reldatabase_id, " +
" tbl.reldatabase, " +
" tbl.rel_createsql, " +
" tbl.account_id, " +
" tbl.`constraint` " +
"FROM `mo_catalog`.`mo_tables` tbl " +
"WHERE " +
" account_id IN (%s) " +
" tbl.account_id IN (%s) " +
"%s" +
"%s" +
" AND relkind = '%s' " +
" AND reldatabase NOT IN (%s)"
" AND tbl.relkind = '%s' " +
" AND tbl.reldatabase NOT IN (%s)"
CDCInsertMOISCPLogSqlTemplate = `REPLACE INTO mo_catalog.mo_iscp_log (` +
`account_id,` +
`table_id,` +
Expand Down Expand Up @@ -475,6 +476,7 @@ var CDCSQLTemplates = [CDCSqlTemplateCount]struct {
"reldatabase",
"rel_createsql",
"account_id",
"constraint",
},
},
CDCGetWatermarkWhereSqlTemplate_Idx: {
Expand Down Expand Up @@ -1071,13 +1073,13 @@ func (b cdcSQLBuilder) CollectTableInfoSQL(accountIDs string, dbNames string, ta
if dbNames == "*" {
return ""
}
return " AND reldatabase IN (" + dbNames + ") "
return " AND tbl.reldatabase IN (" + dbNames + ") "
}(),
func() string {
if tableNames == "*" {
return ""
}
return " AND relname IN (" + tableNames + ") "
return " AND tbl.relname IN (" + tableNames + ") "
}(),
catalog.SystemOrdinaryRel,
AddSingleQuotesJoin(catalog.SystemDatabases),
Expand Down
7 changes: 4 additions & 3 deletions pkg/cdc/table_change_stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -587,6 +587,10 @@ func (s *TableChangeStream) cleanup(ctx context.Context) {
)
defer s.wg.Done()
defer func() {
// Keep ownership until cleanup finishes so a replacement reader cannot
// start while this stream is still closing its sinker/watermark state.
s.runningReaders.CompareAndDelete(s.runningReaderKey, s)

// Decrement table stream state gauge on cleanup
if s.progressTracker != nil {
state, _ := s.progressTracker.GetState()
Expand All @@ -603,9 +607,6 @@ func (s *TableChangeStream) cleanup(ctx context.Context) {
)
}()

// Remove from running readers
s.runningReaders.Delete(s.runningReaderKey)

// Remove watermark cache
removeStart := time.Now()
if err := s.watermarkUpdater.RemoveCachedWM(ctx, s.watermarkKey); err != nil {
Expand Down
106 changes: 106 additions & 0 deletions pkg/cdc/table_change_stream_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -378,6 +378,90 @@ func TestTableChangeStream_Run_DuplicateReader(t *testing.T) {
}
}

func TestTableChangeStreamCleanupDoesNotDeleteReplacedReader(t *testing.T) {
runningReaders := &sync.Map{}

key := "db1.t1"
oldStream := &TableChangeStream{
accountId: 1,
taskId: "task1",
tableInfo: &DbTableInfo{SourceDbName: "db1", SourceTblName: "t1", SourceTblId: 1},
sinker: newTableStreamRecordingSinker(),
watermarkUpdater: newWatermarkUpdaterStub(),
watermarkKey: &WatermarkKey{AccountId: 1, TaskId: "task1", DBName: "db1", TableName: "t1"},
runningReaders: runningReaders,
runningReaderKey: key,
progressTracker: nil,
watermarkStallThreshold: defaultWatermarkStallThreshold,
}
newStream := &TableChangeStream{
accountId: 1,
taskId: "task1",
tableInfo: &DbTableInfo{SourceDbName: "db1", SourceTblName: "t1", SourceTblId: 1},
sinker: newTableStreamRecordingSinker(),
watermarkUpdater: newWatermarkUpdaterStub(),
watermarkKey: &WatermarkKey{AccountId: 1, TaskId: "task1", DBName: "db1", TableName: "t1"},
runningReaders: runningReaders,
runningReaderKey: key,
}
runningReaders.Store(key, oldStream)
runningReaders.Store(key, newStream)

oldStream.wg.Add(1)
oldStream.cleanup(context.Background())

stored, ok := runningReaders.Load(key)
require.True(t, ok, "replaced reader ownership should remain")
require.Same(t, newStream, stored, "old cleanup must not delete a newer reader")
}

func TestTableChangeStreamCleanupKeepsOwnershipUntilCloseFinishes(t *testing.T) {
runningReaders := &sync.Map{}

key := "db1.t1"
sinker := newBlockingCloseSinker()
stream := &TableChangeStream{
accountId: 1,
taskId: "task1",
tableInfo: &DbTableInfo{SourceDbName: "db1", SourceTblName: "t1", SourceTblId: 1},
sinker: sinker,
watermarkUpdater: newWatermarkUpdaterStub(),
watermarkKey: &WatermarkKey{AccountId: 1, TaskId: "task1", DBName: "db1", TableName: "t1"},
runningReaders: runningReaders,
runningReaderKey: key,
progressTracker: nil,
watermarkStallThreshold: defaultWatermarkStallThreshold,
}
runningReaders.Store(key, stream)

stream.wg.Add(1)
cleanupDone := make(chan struct{})
go func() {
stream.cleanup(context.Background())
close(cleanupDone)
}()

select {
case <-sinker.closeStarted:
case <-time.After(time.Second):
t.Fatal("expected cleanup to reach sinker close")
}

stored, ok := runningReaders.Load(key)
require.True(t, ok, "reader ownership should remain while cleanup is still closing")
require.Same(t, stream, stored, "old stream should keep ownership until cleanup finishes")

close(sinker.unblockClose)
select {
case <-cleanupDone:
case <-time.After(time.Second):
t.Fatal("cleanup did not finish after unblocking close")
}

_, ok = runningReaders.Load(key)
require.False(t, ok, "reader ownership should be removed after cleanup finishes")
}

// Integration: commit failure triggers EnsureCleanup rollback, then recovery succeeds
func TestTableChangeStream_CommitFailure_EnsureCleanup_ThenRecover(t *testing.T) {
updaterStub := newWatermarkUpdaterStub()
Expand Down Expand Up @@ -2453,6 +2537,28 @@ func newTableStreamRecordingSinker() *tableStreamRecordingSinker {
return &tableStreamRecordingSinker{recordingSinker: newRecordingSinker()}
}

type blockingCloseSinker struct {
*tableStreamRecordingSinker
closeStarted chan struct{}
unblockClose chan struct{}
closeOnce sync.Once
}

func newBlockingCloseSinker() *blockingCloseSinker {
return &blockingCloseSinker{
tableStreamRecordingSinker: newTableStreamRecordingSinker(),
closeStarted: make(chan struct{}),
unblockClose: make(chan struct{}),
}
}

func (s *blockingCloseSinker) Close() {
s.closeOnce.Do(func() {
close(s.closeStarted)
})
<-s.unblockClose
}

func (s *tableStreamRecordingSinker) Sink(ctx context.Context, data *DecoderOutput) {
s.record("sink")
s.mu.Lock()
Expand Down
58 changes: 49 additions & 9 deletions pkg/cdc/table_scanner.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@ import (
"fmt"
"runtime/debug"
"slices"
"strings"
"sync"
"sync/atomic"
"time"
Expand All @@ -37,6 +36,7 @@ import (
"github.com/matrixorigin/matrixone/pkg/defines"
"github.com/matrixorigin/matrixone/pkg/logutil"
"github.com/matrixorigin/matrixone/pkg/util/executor"
"github.com/matrixorigin/matrixone/pkg/vm/engine"
)

const (
Expand Down Expand Up @@ -744,6 +744,7 @@ func (s *TableDetector) scanTable() error {
}
defer result.Close()

var scanErr error
result.ReadRows(func(rows int, cols []*vector.Vector) bool {
for i := 0; i < rows; i++ {
tblId := vector.MustFixedColWithTypeCheck[uint64](cols[0])[i]
Expand All @@ -752,9 +753,21 @@ func (s *TableDetector) scanTable() error {
dbName := cols[3].GetStringAt(i)
createSql := cols[4].GetStringAt(i)
accountId := vector.MustFixedColWithTypeCheck[uint32](cols[5])[i]
hasForeignKey, decodeErr := tableHasForeignKeyConstraint(cols[6].GetBytesAt(i))
if decodeErr != nil {
scanErr = decodeErr
logutil.Warn(
"cdc.table_detector.scan_constraint_failed",
zap.Uint32("account-id", accountId),
zap.String("db", dbName),
zap.String("table", tblName),
zap.Error(decodeErr),
)
return false
}

// skip table with foreign key
if strings.Contains(strings.ToLower("createSql"), "foreign key") {
if hasForeignKey {
continue
}

Expand All @@ -776,21 +789,48 @@ func (s *TableDetector) scanTable() error {
mp[accountId][key] = newInfo
} else {
idChanged := oldInfo.OnlyDiffinTblId(newInfo)
oldInfo.SourceDbId = dbId
oldInfo.SourceDbName = dbName
oldInfo.SourceTblId = tblId
oldInfo.SourceTblName = tblName
oldInfo.SourceCreateSql = createSql
oldInfo.IdChanged = oldInfo.IdChanged || idChanged
mp[accountId][key] = oldInfo
updatedInfo := oldInfo.Clone()
updatedInfo.SourceDbId = dbId
updatedInfo.SourceDbName = dbName
updatedInfo.SourceTblId = tblId
updatedInfo.SourceTblName = tblName
updatedInfo.SourceCreateSql = createSql
updatedInfo.IdChanged = updatedInfo.IdChanged || idChanged
mp[accountId][key] = updatedInfo
}
}
return true
})
if scanErr != nil {
return scanErr
}

// replace the old table map
s.mu.Lock()
s.Mp = mp
s.mu.Unlock()
return nil
}

func tableHasForeignKeyConstraint(data []byte) (hasForeignKey bool, err error) {
if len(data) == 0 {
return false, nil
}

defer func() {
if r := recover(); r != nil {
err = moerr.NewInternalErrorNoCtxf("unmarshal table constraint failed: %v", r)
}
}()

constraintDef := &engine.ConstraintDef{}
if err := constraintDef.UnmarshalBinary(data); err != nil {
return false, err
}
for _, constraint := range constraintDef.Cts {
if foreignKeyDef, ok := constraint.(*engine.ForeignKeyDef); ok && len(foreignKeyDef.Fkeys) > 0 {
return true, nil
}
}
return false, nil
}
Loading
Loading