Skip to content

Commit 43e0ece

Browse files
philcunliffeclaude
andauthored
Support format-version 3 in icebergRewrite (preserve row lineage) (#24)
* Support format-version 3 in icebergRewrite (preserve row lineage) A v3 rewrite must not renumber existing rows: each surviving row's _row_id and _last_updated_sequence_number are now materialized as explicit columns (reserved field ids 2147483540 / 2147483539) in the rewritten parquet files, since the global sort destroys the positional ordering that derivation relies on. When every live row carries lineage, the new manifest's first_row_id is pinned to the minimum carried _row_id so assignFirstRowIds consumes no new ids and next-row-id does not advance. Tables upgraded from v2 whose rows were never assigned ids fall back to spec-default commit-time assignment. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix typecheck: narrow formatVersion to 2 | 3 in rewrite test helper Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
1 parent 874fa7d commit 43e0ece

2 files changed

Lines changed: 209 additions & 15 deletions

File tree

src/write/rewrite.js

Lines changed: 51 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ import { computeColumnStats } from './stats.js'
1414
import { checkWriteFormat, newSnapshotId, resolveParquetCodec } from './stage.js'
1515

1616
/**
17-
* @import {Manifest, Resolver, Snapshot, StagedUpdate, TableMetadata} from '../../src/types.js'
17+
* @import {Manifest, Resolver, Schema, Snapshot, StagedUpdate, TableMetadata} from '../../src/types.js'
1818
*/
1919

2020
/**
@@ -30,8 +30,16 @@ import { checkWriteFormat, newSnapshotId, resolveParquetCodec } from './stage.js
3030
* boundary). Row contents and counts are preserved (modulo deleted rows and
3131
* order); deletes are consumed, so the new snapshot has no delete files.
3232
*
33-
* v2 only for now: a v3 rewrite would have to preserve `_row_id` row lineage
34-
* rather than let `assignFirstRowIds` renumber the rewritten rows.
33+
* On v3 tables, row lineage is preserved: each surviving row's `_row_id` and
34+
* `_last_updated_sequence_number` are materialized as explicit columns in the
35+
* rewritten files (the global sort breaks positional derivation, so stored
36+
* values are required). When every live row carries lineage, the new manifest's
37+
* `first_row_id` is pinned to the minimum carried `_row_id` so no new row ids
38+
* are consumed (`next-row-id` does not advance). When some rows lack lineage
39+
* (a table upgraded from v2 whose pre-upgrade rows were never assigned ids),
40+
* the manifest is left for commit-time assignment per the spec: stored ids
41+
* still win on read, null rows get derived ids, and `next-row-id` advances by
42+
* the manifest's row count.
3543
*
3644
* @param {object} options
3745
* @param {string} options.tableUrl
@@ -48,10 +56,10 @@ export async function icebergStageRewrite({
4856
if (!tableUrl) throw new Error('tableUrl is required')
4957
if (!resolver?.writer) throw new Error('resolver.writer is required')
5058
const writerFn = resolver.writer
51-
const formatVersion = metadata['format-version']
52-
if (formatVersion !== 2) {
53-
throw new Error(`icebergRewrite supports format-version 2 only (got ${formatVersion}); v3 row lineage is not yet handled`)
59+
if (metadata['format-version'] !== 2 && metadata['format-version'] !== 3) {
60+
throw new Error(`unsupported format-version: ${metadata['format-version']}`)
5461
}
62+
const formatVersion = /** @type {2|3} */ (metadata['format-version'])
5563
if (targetFileRows !== undefined && !(targetFileRows > 0)) {
5664
throw new Error('targetFileRows must be a positive number')
5765
}
@@ -77,10 +85,21 @@ export async function icebergStageRewrite({
7785
checkWriteFormat(metadata.properties?.['write.format.default'])
7886
const codec = resolveParquetCodec(metadata.properties?.['write.parquet.compression-codec'])
7987

80-
// Read every live row (deletes applied), then sort globally.
88+
// Read every live row (deletes applied), then sort globally. For v3 tables
89+
// the rows carry `_row_id` / `_last_updated_sequence_number` (derived or
90+
// stored by the read path).
8191
const liveRows = await icebergRead({ tableUrl, metadata, resolver })
8292
const sortedRows = comparator ? [...liveRows].sort(comparator) : liveRows
8393

94+
const rowLineage = formatVersion >= 3
95+
// Preserve mode requires complete lineage: rows from pre-upgrade v2
96+
// snapshots read with null ids and need commit-time assignment instead.
97+
const allLineage = rowLineage && liveRows.length > 0 &&
98+
liveRows.every(r => r._row_id != null && r._last_updated_sequence_number != null)
99+
const minRowId = allLineage
100+
? liveRows.reduce((min, r) => r._row_id < min ? r._row_id : min, liveRows[0]._row_id)
101+
: undefined
102+
84103
// Regroup under the target partition spec (re-derives tuples from values, so
85104
// files written under an older spec are rewritten under the new one).
86105
const groups = partitionSpec.fields.length
@@ -90,6 +109,22 @@ export async function icebergStageRewrite({
90109
const snapshotId = newSnapshotId(metadata)
91110
const manifestUuid = uuid4()
92111

112+
// For v3, materialize the carried lineage as explicit columns in the
113+
// rewritten files (reserved field ids per spec). The extended schema is
114+
// passed to the parquet writer only: stats, partitioning, and the manifest's
115+
// embedded schema stay user-only.
116+
/** @type {Schema} */
117+
const writeSchema = rowLineage
118+
? {
119+
...schema,
120+
fields: [
121+
...schema.fields,
122+
{ id: 2147483540, name: '_row_id', required: false, type: 'long' },
123+
{ id: 2147483539, name: '_last_updated_sequence_number', required: false, type: 'long' },
124+
],
125+
}
126+
: schema
127+
93128
/** @type {{ partition: Record<string, any>, dataFile: any, path: string }[]} */
94129
const writtenDataFiles = []
95130
for (const group of groups) {
@@ -98,7 +133,7 @@ export async function icebergStageRewrite({
98133
if (chunk.length === 0) continue
99134
const dataPath = `${tableUrl}/data/${uuid4()}.parquet`
100135
const dataWriter = writerFn(dataPath)
101-
await writeParquet({ writer: dataWriter, schema, records: chunk, codec })
136+
await writeParquet({ writer: dataWriter, schema: writeSchema, records: chunk, codec })
102137
const stats = computeColumnStats(chunk, schema)
103138
writtenDataFiles.push({
104139
partition: group.partition,
@@ -160,6 +195,14 @@ export async function icebergStageRewrite({
160195
deleted_rows_count: 0n,
161196
partitions,
162197
}
198+
if (allLineage) {
199+
// Every rewritten row carries a materialized `_row_id`, so the manifest
200+
// does not need a fresh id range from `assignFirstRowIds`. Pin its
201+
// first_row_id to the smallest carried id: ids are distinct and below
202+
// next-row-id, so min + row count never exceeds next-row-id and no new
203+
// ids are consumed (`next-row-id` does not advance for a pure rewrite).
204+
newManifest.first_row_id = minRowId
205+
}
163206

164207
// Supersede every prior manifest (data + delete) for the rewritten snapshot.
165208
const priorManifests = await loadPriorManifests(metadata, resolver)

test/write/rewrite.test.js

Lines changed: 158 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,13 @@
11
import { describe, expect, it, vi } from 'vitest'
22
import { fileCatalog } from '../../src/catalog/file.js'
3+
import { fetchAvroRecords } from '../../src/fetch.js'
34
import { fileCatalogCommit } from '../../src/write/commit.js'
45
import { icebergCreate } from '../../src/create.js'
56
import { icebergManifests, splitManifestEntries } from '../../src/manifest.js'
67
import { icebergRead } from '../../src/read.js'
78
import { icebergAppend, icebergRewrite } from '../../src/write/write.js'
89
import { icebergStageAppend } from '../../src/write/stage.js'
10+
import { icebergStageDeletionVector } from '../../src/write/stage-deletion-vector.js'
911
import { icebergStagePositionDelete } from '../../src/write/stage-position-delete.js'
1012
import { icebergStageRewrite } from '../../src/write/rewrite.js'
1113
import { deserializeValue } from '../../src/write/serde.js'
@@ -54,17 +56,27 @@ function longBound(entry, fieldId, side) {
5456
return bytes ? deserializeValue(bytes, 'long') : undefined
5557
}
5658

59+
/**
60+
* Map each row's id to its lineage pair, for before/after comparison.
61+
* @param {Record<string, any>[]} rows
62+
* @returns {Map<any, {rowId: any, lusn: any}>}
63+
*/
64+
function lineageById(rows) {
65+
return new Map(rows.map(r => [r.id, { rowId: r._row_id, lusn: r._last_updated_sequence_number }]))
66+
}
67+
5768
/**
5869
* Create a sorted, multi-file table and return its committed metadata.
5970
* @param {object} [opts]
6071
* @param {SortOrder} [opts.sortOrder]
72+
* @param {2 | 3} [opts.formatVersion]
6173
* @returns {Promise<{ tableUrl: string, resolver: import('../../src/types.js').Resolver, metadata: TableMetadata }>}
6274
*/
63-
async function makeMultiFileTable({ sortOrder } = {}) {
75+
async function makeMultiFileTable({ sortOrder, formatVersion } = {}) {
6476
vi.spyOn(Date, 'now').mockReturnValue(1700000000000)
6577
const tableUrl = 'mem://rewrite'
6678
const { resolver } = memResolver()
67-
let metadata = await icebergCreate({ tableUrl, resolver, schema, sortOrder })
79+
let metadata = await icebergCreate({ tableUrl, resolver, schema, sortOrder, formatVersion })
6880
const batches = [
6981
[{ id: 5n, name: 'e' }, { id: 2n, name: 'b' }],
7082
[{ id: 1n, name: 'a' }, { id: 6n, name: 'f' }],
@@ -209,16 +221,155 @@ describe('icebergRewrite — safety', () => {
209221
expect(rows.map(r => r.id).sort()).toEqual([1n, 2n])
210222
})
211223

212-
it('rejects format-version 3 (row lineage not yet handled)', async () => {
224+
it('rejects unsupported format versions', async () => {
225+
const { tableUrl, resolver, metadata } = await makeMultiFileTable()
226+
await expect(() => icebergStageRewrite({
227+
tableUrl, metadata: { ...metadata, 'format-version': 1 }, resolver,
228+
})).rejects.toThrow(/unsupported format-version: 1/)
229+
})
230+
})
231+
232+
describe('icebergRewrite — v3 row lineage', () => {
233+
it('preserves _row_id and _last_updated_sequence_number across a sorted rewrite', async () => {
234+
const { tableUrl, resolver, metadata } = await makeMultiFileTable({ sortOrder: sortById, formatVersion: 3 })
235+
expect(metadata['next-row-id']).toBe(6)
236+
const before = await icebergRead({ tableUrl, metadata, resolver })
237+
expect(new Set(before.map(r => r._row_id))).toEqual(new Set([0n, 1n, 2n, 3n, 4n, 5n]))
238+
239+
const staged = await icebergStageRewrite({ tableUrl, metadata, resolver })
240+
// A pure rewrite carries existing rows, so it assigns no new row ids.
241+
expect(staged.snapshot['first-row-id']).toBe(6)
242+
expect(staged.snapshot['added-rows']).toBe(0)
243+
expect(staged.requirements).toContainEqual({ type: 'assert-next-row-id', 'next-row-id': 6 })
244+
245+
const after = await fileCatalogCommit({ tableUrl, metadata, staged, resolver })
246+
expect(after['next-row-id']).toBe(6)
247+
248+
const rows = await icebergRead({ tableUrl, metadata: after, resolver })
249+
expect(rows.map(r => r.id)).toEqual([1n, 2n, 3n, 4n, 5n, 6n])
250+
// The global sort scrambled row positions, so preserved values must come
251+
// from the materialized _row_id column, not positional derivation.
252+
expect(rows.map(r => r._row_id)).not.toEqual([0n, 1n, 2n, 3n, 4n, 5n])
253+
expect(lineageById(rows)).toEqual(lineageById(before))
254+
255+
// The rewritten manifest is pinned to the carried range: min(_row_id).
256+
const snapshot = after.snapshots?.find(s => s['snapshot-id'] === after['current-snapshot-id'])
257+
const manifests = await fetchAvroRecords(snapshot?.['manifest-list'] ?? '', resolver)
258+
expect(manifests.length).toBe(1)
259+
expect(manifests[0].first_row_id).toBe(0n)
260+
})
261+
262+
it('preserves lineage across split output files', async () => {
263+
const { tableUrl, resolver, metadata } = await makeMultiFileTable({ sortOrder: sortById, formatVersion: 3 })
264+
const before = await icebergRead({ tableUrl, metadata, resolver })
265+
266+
const staged = await icebergStageRewrite({ tableUrl, metadata, resolver, targetFileRows: 2 })
267+
const after = await fileCatalogCommit({ tableUrl, metadata, staged, resolver })
268+
269+
const entries = splitManifestEntries(await icebergManifests({ metadata: after, resolver })).dataEntries
270+
expect(entries.length).toBe(3)
271+
const rows = await icebergRead({ tableUrl, metadata: after, resolver })
272+
expect(lineageById(rows)).toEqual(lineageById(before))
273+
expect(after['next-row-id']).toBe(6)
274+
})
275+
276+
it('consumes deletion vectors and keeps surviving rows\' lineage', async () => {
213277
vi.spyOn(Date, 'now').mockReturnValue(1700000000000)
214-
const tableUrl = 'mem://rewrite-v3'
278+
const tableUrl = 'mem://rewrite-v3-dv'
215279
const { resolver } = memResolver()
216280
const created = await icebergCreate({ tableUrl, resolver, schema, formatVersion: 3 })
217-
const appended = await icebergStageAppend({ tableUrl, metadata: created, records: [{ id: 1n, name: 'a' }], resolver })
281+
const records = [{ id: 1n, name: 'a' }, { id: 2n, name: 'b' }, { id: 3n, name: 'c' }]
282+
const appended = await icebergStageAppend({ tableUrl, metadata: created, records, resolver })
218283
const afterAppend = await fileCatalogCommit({ tableUrl, metadata: created, staged: appended, resolver })
284+
const dataPath = appended.writtenFiles[0]
219285

220-
await expect(() => icebergStageRewrite({ tableUrl, metadata: afterAppend, resolver }))
221-
.rejects.toThrow(/format-version 2 only/)
286+
const delStaged = await icebergStageDeletionVector({
287+
tableUrl, metadata: afterAppend, deletes: [{ file_path: dataPath, pos: 1n }], resolver,
288+
})
289+
const afterDelete = await fileCatalogCommit({ tableUrl, metadata: afterAppend, staged: delStaged, resolver })
290+
291+
const staged = await icebergStageRewrite({ tableUrl, metadata: afterDelete, resolver })
292+
const afterRewrite = await fileCatalogCommit({ tableUrl, metadata: afterDelete, staged, resolver })
293+
294+
const rows = await icebergRead({ tableUrl, metadata: afterRewrite, resolver })
295+
expect(rows.map(r => ({ id: r.id, _row_id: r._row_id, _last_updated_sequence_number: r._last_updated_sequence_number }))).toEqual([
296+
{ id: 1n, _row_id: 0n, _last_updated_sequence_number: 1n },
297+
{ id: 3n, _row_id: 2n, _last_updated_sequence_number: 1n },
298+
])
299+
const { dataEntries, deleteEntries } = splitManifestEntries(await icebergManifests({ metadata: afterRewrite, resolver }))
300+
expect(deleteEntries.length).toBe(0)
301+
expect(dataEntries.length).toBe(1)
302+
// The deleted row's id (1n) is retired with it; next-row-id is unchanged.
303+
expect(afterRewrite['next-row-id']).toBe(3)
304+
})
305+
306+
it('appends after a rewrite continue from the unchanged next-row-id', async () => {
307+
const { tableUrl, resolver, metadata } = await makeMultiFileTable({ sortOrder: sortById, formatVersion: 3 })
308+
const staged = await icebergStageRewrite({ tableUrl, metadata, resolver })
309+
const after = await fileCatalogCommit({ tableUrl, metadata, staged, resolver })
310+
311+
const appended = await icebergStageAppend({
312+
tableUrl, metadata: after, records: [{ id: 7n, name: 'g' }, { id: 8n, name: 'h' }], resolver,
313+
})
314+
const final = await fileCatalogCommit({ tableUrl, metadata: after, staged: appended, resolver })
315+
expect(final['next-row-id']).toBe(8)
316+
317+
const rows = await icebergRead({ tableUrl, metadata: final, resolver })
318+
// New rows take fresh ids 6n and 7n; no collision with the preserved 0n-5n.
319+
expect(new Set(rows.map(r => r._row_id))).toEqual(new Set([0n, 1n, 2n, 3n, 4n, 5n, 6n, 7n]))
320+
expect(rows.find(r => r.id === 7n)?._row_id).toBe(6n)
321+
expect(rows.find(r => r.id === 8n)?._row_id).toBe(7n)
322+
})
323+
324+
it('keeps lineage stable through a rewrite of a rewrite', async () => {
325+
const { tableUrl, resolver, metadata } = await makeMultiFileTable({ sortOrder: sortById, formatVersion: 3 })
326+
const before = await icebergRead({ tableUrl, metadata, resolver })
327+
328+
const staged1 = await icebergStageRewrite({ tableUrl, metadata, resolver })
329+
const after1 = await fileCatalogCommit({ tableUrl, metadata, staged: staged1, resolver })
330+
// Source files now carry materialized lineage columns.
331+
const staged2 = await icebergStageRewrite({ tableUrl, metadata: after1, resolver, targetFileRows: 3 })
332+
const after2 = await fileCatalogCommit({ tableUrl, metadata: after1, staged: staged2, resolver })
333+
334+
const rows = await icebergRead({ tableUrl, metadata: after2, resolver })
335+
expect(lineageById(rows)).toEqual(lineageById(before))
336+
expect(after2['next-row-id']).toBe(6)
337+
})
338+
339+
it('assigns fresh row ids when rewriting an upgraded table whose rows lack lineage', async () => {
340+
vi.spyOn(Date, 'now').mockReturnValue(1700000000000)
341+
const tableUrl = 'mem://rewrite-v3-upgrade'
342+
const { resolver } = memResolver()
343+
const created = await icebergCreate({ tableUrl, resolver, schema })
344+
const records = [{ id: 1n, name: 'a' }, { id: 2n, name: 'b' }, { id: 3n, name: 'c' }]
345+
const appended = await icebergStageAppend({ tableUrl, metadata: created, records, resolver })
346+
const afterAppend = await fileCatalogCommit({ tableUrl, metadata: created, staged: appended, resolver })
347+
/** @type {TableMetadata} */
348+
const upgraded = { ...afterAppend, 'format-version': 3, 'next-row-id': 0 }
349+
350+
// Pre-upgrade rows read with null lineage.
351+
const before = await icebergRead({ tableUrl, metadata: upgraded, resolver })
352+
expect(before.every(r => r._row_id === null)).toBe(true)
353+
354+
const staged = await icebergStageRewrite({ tableUrl, metadata: upgraded, resolver })
355+
// No carried ids, so the rewritten manifest gets a fresh range at commit.
356+
expect(staged.snapshot['first-row-id']).toBe(0)
357+
expect(staged.snapshot['added-rows']).toBe(3)
358+
359+
const after = await fileCatalogCommit({ tableUrl, metadata: upgraded, staged, resolver })
360+
expect(after['next-row-id']).toBe(3)
361+
const rows = await icebergRead({ tableUrl, metadata: after, resolver })
362+
expect(rows.map(r => r._row_id)).toEqual([0n, 1n, 2n])
363+
})
364+
365+
it('detects a concurrent v3 commit via assert-next-row-id', async () => {
366+
const { tableUrl, resolver, metadata } = await makeMultiFileTable({ sortOrder: sortById, formatVersion: 3 })
367+
const staged = await icebergStageRewrite({ tableUrl, metadata, resolver })
368+
// Committing against metadata whose next-row-id moved (a concurrent append
369+
// that somehow passed the ref check) must fail rather than risk id reuse.
370+
await expect(fileCatalogCommit({
371+
tableUrl, metadata: { ...metadata, 'next-row-id': 8 }, staged, resolver,
372+
})).rejects.toThrow(/next-row-id expected 6, got 8/)
222373
})
223374
})
224375

0 commit comments

Comments
 (0)