Skip to content

Commit 8480da0

Browse files
fix cdc watermark add failure handling (#25925)
- Updated pkg/cdc/watermark_updater.go: execAddWM now returns immediately when the catalog INSERT fails, so it no longer writes the watermark into cacheCommitted. The existing onJobs cleanup path then propagates the error to the waiting job. - Updated pkg/cdc/watermark_updater_test.go: added TestAuditAddWatermarkFailureIsReturnedAndNotCached, covering the case where SELECT returns no watermark row, INSERT fails with an injected error, the job receives that error, and the failed watermark is not cached. Approved by: @gouhongshen
1 parent 3046f14 commit 8480da0

2 files changed

Lines changed: 54 additions & 0 deletions

File tree

pkg/cdc/watermark_updater.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -813,6 +813,7 @@ func (u *CDCWatermarkUpdater) execAddWM() (errMsg string, err error) {
813813
err = u.ie.Exec(ctx, addSql, ie.SessionOverrideOptions{})
814814
if err != nil {
815815
errMsg = fmt.Sprintf("add sql \"%s\" failed", addSql)
816+
return
816817
}
817818
u.Lock()
818819
defer u.Unlock()

pkg/cdc/watermark_updater_test.go

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,10 @@ type retryableMockExecutor struct {
4848
lastSQL string
4949
}
5050

51+
type failAddWatermarkExecutor struct {
52+
insertErr error
53+
}
54+
5155
func (m *retryableMockExecutor) Exec(_ context.Context, sql string, _ ie.SessionOverrideOptions) error {
5256
m.mu.Lock()
5357
defer m.mu.Unlock()
@@ -66,6 +70,30 @@ func (m *retryableMockExecutor) Query(_ context.Context, _ string, _ ie.SessionO
6670

6771
func (m *retryableMockExecutor) ApplySessionOverride(_ ie.SessionOverrideOptions) {}
6872

73+
func (m *failAddWatermarkExecutor) Exec(_ context.Context, sql string, _ ie.SessionOverrideOptions) error {
74+
if strings.HasPrefix(sql, "INSERT INTO `mo_catalog`.`mo_cdc_watermark`") {
75+
return m.insertErr
76+
}
77+
return nil
78+
}
79+
80+
func (m *failAddWatermarkExecutor) Query(_ context.Context, sql string, _ ie.SessionOverrideOptions) ie.InternalExecResult {
81+
if strings.HasPrefix(sql, "SELECT") {
82+
return &InternalExecResultForTest{
83+
resultSet: &MysqlResultSetForTest{
84+
Data: [][]interface{}{},
85+
},
86+
}
87+
}
88+
return &InternalExecResultForTest{
89+
resultSet: &MysqlResultSetForTest{
90+
Data: [][]interface{}{},
91+
},
92+
}
93+
}
94+
95+
func (m *failAddWatermarkExecutor) ApplySessionOverride(_ ie.SessionOverrideOptions) {}
96+
6997
func newWmMockSQLExecutor() *wmMockSQLExecutor {
7098
return &wmMockSQLExecutor{
7199
mp: make(map[string]string),
@@ -116,6 +144,31 @@ func (m *wmMockSQLExecutor) Query(ctx context.Context, sql string, pts ie.Sessio
116144

117145
func (m *wmMockSQLExecutor) ApplySessionOverride(opts ie.SessionOverrideOptions) {}
118146

147+
func TestAuditAddWatermarkFailureIsReturnedAndNotCached(t *testing.T) {
148+
insertErr := moerr.NewInternalErrorNoCtx("injected watermark insert failure")
149+
updater := NewCDCWatermarkUpdater("add-failure", &failAddWatermarkExecutor{
150+
insertErr: insertErr,
151+
})
152+
key := WatermarkKey{
153+
AccountId: 1,
154+
TaskId: "task",
155+
DBName: "db",
156+
TableName: "tbl",
157+
}
158+
watermark := types.BuildTS(10, 1)
159+
job := NewGetOrAddCommittedWMJob(context.Background(), &key, &watermark)
160+
161+
updater.onJobs(job)
162+
163+
result := job.GetResult()
164+
require.ErrorIs(t, result.Err, insertErr)
165+
166+
updater.RLock()
167+
_, ok := updater.cacheCommitted[key]
168+
updater.RUnlock()
169+
require.False(t, ok)
170+
}
171+
119172
func TestWatermarkUpdater_CommitRetrySuccess(t *testing.T) {
120173
exec := &retryableMockExecutor{failRemaining: 1}
121174
updater := NewCDCWatermarkUpdater("retry-success", exec)

0 commit comments

Comments
 (0)