Skip to content
Merged
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
148 changes: 117 additions & 31 deletions packages/plugin-backfill/src/mv-replay-plan.e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

// ---------------------------------------------------------------------------
Expand All @@ -36,14 +38,22 @@ 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
// The execute loop submits + polls in parallel; a stateless executor avoids
// 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
Expand All @@ -61,6 +71,65 @@ async function aggregateByBucket(fqn: string, valueExpr: string): Promise<Array<
)
}

// Chunk queries are fire-and-forget on a stateless client, so each one can
// land on a different replica than the session that inserted the source.
// select_sequential_consistency does not fetch parts that were never quorum
// commits; it only hides or rejects blocks the quorum has not confirmed. A
// replica that has not attached the partition therefore finishes the INSERT
// with 0 written rows. Sample a fresh session per replica, and on
// replicated/shared engines sync that session's replica before counting.
async function waitUntilSourceVisible(): Promise<void> {
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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Require distinct replicas before declaring the source ready

When the load balancer routes multiple fresh sessions to the same healthy replica, each successful count increments readySamples, so this can reach active_replicas without ever checking a lagging replica. A subsequent stateless chunk can still land on that unchecked replica and reproduce the 0-row replay this change is intended to prevent. Record a replica identity with each sample and require the expected number of distinct ready replicas, or target every replica explicitly.

Useful? React with 👍 / 👎.

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<string, number | undefined> {
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.
Expand All @@ -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 <T>(
sql: string,
settings?: Record<string, string | number | boolean | undefined>,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)')
Expand All @@ -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)
})
Loading