diff --git a/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts b/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts index 2d3b5f5c..83e7c7f2 100644 --- a/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts +++ b/packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts @@ -11,15 +11,17 @@ import { createStatelessLiveExecutor, getRequiredEnv, pollUntil, + quoteIdent, waitForTable, + type LiveEnv, } from '@chkit/clickhouse/e2e-testkit' -import { executeBackfill } from './async-backfill.js' +import { executeBackfill, type BackfillResult } from './async-backfill.js' import { buildChunkExecutionSql } from './chunking/sql.js' import { generateIdempotencyToken } from './chunking/utils/ids.js' import { PlanSchema } from './options.js' import { buildBackfillPlan } from './planner.js' -import type { PlannerQuery } from './chunking/types.js' +import type { Chunk, PlannerQuery } from './chunking/types.js' import type { BackfillPlanState } from './types.js' // --------------------------------------------------------------------------- @@ -36,6 +38,9 @@ import type { BackfillPlanState } from './types.js' const SOURCE_ROWS = 4000 const BUCKETS = 4 +// Replaying a chunk under the same plan id reuses its dedup token. An empty +// INSERT can record that token, and the next attempt then commits nothing again. +const MAX_EMPTY_CHUNK_REPLAYS = 3 // DDL / inserts / counts go through the session-bound executor (sequential). let ddl: ClickHouseExecutor @@ -43,7 +48,12 @@ let ddl: ClickHouseExecutor // ObsessionDB session-locking errors under concurrency. let runExecutor: ClickHouseExecutor let plannerQuery: PlannerQuery +let liveEnv: LiveEnv let db: string +// Plain MergeTree has one copy of the data. Replicated/Shared engines need +// every active replica to attach the source parts before a chunk reads them. +let syncSourceReplica = false +let sourceReplicas = 1 let sourceTable: string let targetTable: string let sourceFqn: string @@ -61,6 +71,65 @@ async function aggregateByBucket(fqn: string, valueExpr: string): Promise { + let readySamples = 0 + const ready = await pollUntil(async () => { + const session = createLiveExecutor(liveEnv) + try { + if (syncSourceReplica) { + await session.command( + `SYSTEM SYNC REPLICA ${quoteIdent(db)}.${quoteIdent(sourceTable)} LIGHTWEIGHT`, + ) + } + const [row] = await session.query<{ cnt: string }>( + `SELECT toString(count()) AS cnt FROM ${sourceFqn} SETTINGS select_sequential_consistency = 1`, + ) + if (Number(row?.cnt ?? 0) === SOURCE_ROWS) readySamples += 1 + return readySamples + } finally { + await session.close() + } + }, (samples) => samples >= sourceReplicas, { timeoutMs: 45_000, intervalMs: 250 }) + expect(ready, `replicas that can see all ${SOURCE_ROWS} source rows`).toBeGreaterThanOrEqual(sourceReplicas) +} + +function chunkExecutionSql(planId: string, chunk: Chunk, plan: BackfillPlanState): string { + // enable_parallel_replicas is on by default on ObsessionDB. A stale follower + // can answer one partition with an empty scan while the coordinator still + // reports the query finished. Planning already disables it for the same reason. + return `${buildChunkExecutionSql({ + planId, + chunk, + target: plan.target, + sourceTarget: plan.execution.sourceTarget, + table: plan.chunkPlan.table, + mvReplayQueries: plan.execution.mvReplayQueries, + targetColumns: plan.execution.targetColumns, + idempotencyToken: plan.execution.requireIdempotencyToken + ? generateIdempotencyToken(planId, chunk.id) + : '', + })}, select_sequential_consistency = 1, enable_parallel_replicas = 0` +} + +function emptyChunkIds(result: BackfillResult): string[] { + return Object.entries(result.progress) + .filter(([, chunk]) => chunk.status === 'done' && (chunk.writtenRows ?? 0) === 0) + .map(([id]) => id) +} + +function writtenRowsByChunk(result: BackfillResult): Record { + return Object.fromEntries( + Object.entries(result.progress).map(([id, chunk]) => [id, chunk.writtenRows]), + ) +} + function schemaSource(): string { // Plain-object definitions (no imports) so loadSchemaDefinitions can evaluate // the file straight from a temp dir, matching the unit-test convention. @@ -87,10 +156,10 @@ export const events_mv = { } beforeAll(async () => { - const env = getRequiredEnv() - db = env.clickhouseDatabase - ddl = createLiveExecutor(env) - runExecutor = createStatelessLiveExecutor(env) + liveEnv = getRequiredEnv() + db = liveEnv.clickhouseDatabase + ddl = createLiveExecutor(liveEnv) + runExecutor = createStatelessLiveExecutor(liveEnv) plannerQuery = async ( sql: string, settings?: Record, @@ -123,6 +192,19 @@ beforeAll(async () => { await waitForTable(ddl, db, sourceTable) await waitForTable(ddl, db, targetTable) + const [engineRow] = await ddl.query<{ engine: string }>( + `SELECT engine FROM system.tables WHERE database = '${db}' AND name = '${sourceTable}'`, + ) + syncSourceReplica = /Shared|Replicated/.test(engineRow?.engine ?? '') + if (syncSourceReplica) { + const [replicaRow] = await ddl.query<{ active: string }>( + `SELECT toString(active_replicas) AS active + FROM system.replicas + WHERE database = '${db}' AND table = '${sourceTable}'`, + ) + sourceReplicas = Math.max(1, Number(replicaRow?.active ?? 0)) + } + const rows = Array.from({ length: SOURCE_ROWS }, (_, i) => ({ id: i, bucket: i % BUCKETS, @@ -183,35 +265,42 @@ describe('e2e: mv_replay backfill of an empty aggregate target (chkit#187)', () // One chunk per source partition — a real multi-chunk plan over the source. expect(plan.chunkPlan.chunks.length).toBe(BUCKETS) - const result = await executeBackfill({ + const runChunks = (planId: string, chunkIds: string[]) => executeBackfill({ executor: runExecutor, - planId: plan.planId, - chunks: plan.chunkPlan.chunks.map((chunk) => ({ id: chunk.id })), + planId, + chunks: chunkIds.map((id) => ({ id })), buildQuery: ({ id }) => { const planChunk = plan.chunkPlan.chunks.find((candidate) => candidate.id === id) if (!planChunk) throw new Error(`Chunk ${id} not found in plan`) - // The source rows were inserted moments ago and each chunk may run on - // any replica: a lagging one would replay an empty partition and - // "finish" with nothing written. The generated SQL ends in its SETTINGS clause. - return `${buildChunkExecutionSql({ - planId: plan.planId, - chunk: planChunk, - target: plan.target, - sourceTarget: plan.execution.sourceTarget, - table: plan.chunkPlan.table, - mvReplayQueries: plan.execution.mvReplayQueries, - targetColumns: plan.execution.targetColumns, - idempotencyToken: plan.execution.requireIdempotencyToken - ? generateIdempotencyToken(plan.planId, planChunk.id) - : '', - })}, select_sequential_consistency = 1` + return chunkExecutionSql(planId, planChunk, plan) }, concurrency: 3, pollIntervalMs: 1500, }) - expect(result.failed).toBe(0) - expect(result.completed).toBe(plan.chunkPlan.chunks.length) + // The source insert has already returned, but another replica may not + // have attached those parts yet. Sync and count before any chunk reads. + await waitUntilSourceVisible() + + let result = await runChunks(plan.planId, plan.chunkPlan.chunks.map((chunk) => chunk.id)) + const writtenAttempts = [writtenRowsByChunk(result)] + let missing = emptyChunkIds(result) + for (let attempt = 1; missing.length > 0 && attempt <= MAX_EMPTY_CHUNK_REPLAYS; attempt++) { + // A new plan id is a new query id and a new dedup token. Reusing the + // token from the empty INSERT would commit another 0-row replay. + await waitUntilSourceVisible() + const retry = await runChunks(`${plan.planId}-r${attempt}`, missing) + expect(retry.failed).toBe(0) + expect(retry.completed).toBe(missing.length) + writtenAttempts.push(writtenRowsByChunk(retry)) + result = { ...result, progress: { ...result.progress, ...retry.progress } } + missing = emptyChunkIds(retry) + } + + const failed = Object.values(result.progress).filter((chunk) => chunk.status === 'failed').length + const completed = Object.values(result.progress).filter((chunk) => chunk.status === 'done').length + expect(failed).toBe(0) + expect(completed).toBe(plan.chunkPlan.chunks.length) // Per-bucket values must match a forward run of the MV over the whole source. const expected = await aggregateByBucket(sourceFqn, 'sum(id)') @@ -223,9 +312,6 @@ describe('e2e: mv_replay backfill of an empty aggregate target (chkit#187)', () () => aggregateByBucket(targetFqn, 'sum(total)'), (rows) => Bun.deepEquals(rows, expected), ) - const writtenRowsByChunk = Object.fromEntries( - Object.entries(result.progress).map(([id, chunk]) => [id, chunk.writtenRows]), - ) - expect(actual, `rows written per chunk: ${JSON.stringify(writtenRowsByChunk)}`).toEqual(expected) - }, 180_000) + expect(actual, `rows written per chunk attempt: ${JSON.stringify(writtenAttempts)}`).toEqual(expected) + }, 240_000) })