From 354c09be3738445c5e4c44fd288e745acc8e8df1 Mon Sep 17 00:00:00 2001 From: Krishanu Saha Date: Fri, 31 Jul 2026 22:20:38 +0530 Subject: [PATCH] fix(backfill): populate cdc timestamp for snapshot records when cdc enabled --- drivers/abstract/backfill.go | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/drivers/abstract/backfill.go b/drivers/abstract/backfill.go index c488bf4b7..bfedf1613 100644 --- a/drivers/abstract/backfill.go +++ b/drivers/abstract/backfill.go @@ -76,6 +76,8 @@ func (a *AbstractDriver) Backfill(mainCtx context.Context, backfilledStreams cha logger.Infof("Thread[%s]: created writer for chunk min[%s] and max[%s] of stream %s", threadID, chunk.Min, chunk.Max, stream.ID()) + hasCDCTimestamp, _ := stream.GetStream().Schema.GetProperty(constants.CdcTimestamp) + return a.driver.ChunkIterator(backfillCtx, stream, chunk, func(ctx context.Context, data map[string]any, sourceBytes int64) error { olakeID := utils.GetKeysHash(data, stream.GetStream().SourceDefinedPrimaryKey.Array()...) olakeColumns := map[string]any{ @@ -84,8 +86,7 @@ func (a *AbstractDriver) Backfill(mainCtx context.Context, backfilledStreams cha constants.OlakeTimestamp: time.Now().UTC(), } - // Add CDC specific columns only for CDC mode - if stream.GetSyncMode() == types.CDC { + if hasCDCTimestamp { olakeColumns[constants.CdcTimestamp] = time.Unix(0, 0) }