Skip to content
Merged
Show file tree
Hide file tree
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
1 change: 1 addition & 0 deletions src/manifest.js
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ async function fetchManifests(manifests, resolver) {
// Inherit sequence number from manifest if not present in entry
for (const entry of entries) {
entry.partition_spec_id = manifest.partition_spec_id ?? 0
if (entry.snapshot_id == null) entry.snapshot_id = manifest.added_snapshot_id

if (entry.sequence_number === undefined) {
// When reading v1 manifests with no sequence number column,
Expand Down
64 changes: 64 additions & 0 deletions src/write/manifest.js
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,12 @@ function manifestEntrySchema(schema, partitionSpec, formatVersion, manifestConte
mapField('nan_value_counts', 137, 'k138_v139', 138, 139, 'long'),
mapField('lower_bounds', 125, 'k126_v127', 126, 127, 'bytes'),
mapField('upper_bounds', 128, 'k129_v130', 129, 130, 'bytes'),
{
name: 'split_offsets',
type: ['null', { type: 'array', items: 'long', 'element-id': 133 }],
default: null,
'field-id': 132,
},
{ name: 'sort_order_id', type: ['null', 'int'], default: null, 'field-id': 140 },
]
if (manifestContent === 1) {
Expand Down Expand Up @@ -216,6 +222,63 @@ export function writeDeleteManifest({ writer, schema, partitionSpec, snapshotId,
})
}

/**
* Write a data manifest containing already-existing data entries. Used when
* manifests are rewritten without touching the data files they reference —
* merging many small manifests into one, for instance. Carried-over files
* must keep their original data and file sequence numbers rather than
* inheriting the rewriting snapshot's, which would misattribute their data
* sequence numbers and change which delete files apply to them. All entries
* must belong to the supplied partition spec.
*
* @param {object} options
* @param {Writer} options.writer
* @param {Schema} options.schema
* @param {PartitionSpec} options.partitionSpec
* @param {ManifestEntry[]} options.entries
* @param {2|3} [options.formatVersion]
* @returns {void | Promise<void>} resolves when the writer's `finish()` lands
*/
export function writeExistingDataManifest({ writer, schema, partitionSpec, entries, formatVersion = 2 }) {
const records = entries.map(entry => {
const dataFile = entry.data_file
if (entry.status === 2) {
throw new Error('writeExistingDataManifest cannot rewrite deleted entries as existing')
}
if (dataFile.content !== 0) {
throw new Error(`writeExistingDataManifest expects data files (content=0), got content=${dataFile.content}`)
}
if (entry.partition_spec_id !== partitionSpec['spec-id']) {
throw new Error(`existing data entry partition spec ${entry.partition_spec_id} does not match ${partitionSpec['spec-id']}`)
}
const record = manifestEntryRecord(dataFile, schema, partitionSpec, 0n, formatVersion, 0)
record.status = 0
record.snapshot_id = entry.snapshot_id ?? null
record.sequence_number = entry.sequence_number ?? null
record.file_sequence_number = entry.file_sequence_number ?? null
if (record.snapshot_id == null) {
throw new Error('existing data manifest entry missing snapshot id')
}
if (record.sequence_number == null || record.file_sequence_number == null) {
throw new Error('existing data manifest entry missing sequence numbers')
}
return record
})

return avroWrite({
writer,
schema: manifestEntrySchema(schema, partitionSpec, formatVersion, 0),
records,
metadata: {
'format-version': String(formatVersion),
content: 'data',
schema: icebergSchemaJson(schema),
'partition-spec': partitionSpecJson(partitionSpec),
'partition-spec-id': String(partitionSpec['spec-id']),
},
})
}

/**
* Write a delete manifest containing already-existing delete entries. Used
* when a v3 deletion vector replaces an older vector in a mixed manifest: the
Expand Down Expand Up @@ -292,6 +355,7 @@ function manifestEntryRecord(dataFile, schema, partitionSpec, snapshotId, format
nan_value_counts: encodeMap(dataFile.nan_value_counts),
lower_bounds: encodeMap(dataFile.lower_bounds),
upper_bounds: encodeMap(dataFile.upper_bounds),
split_offsets: dataFile.split_offsets?.length ? dataFile.split_offsets : null,
sort_order_id: dataFile.content === 1 ? null : dataFile.sort_order_id ?? 0,
}
if (manifestContent === 1) {
Expand Down
42 changes: 41 additions & 1 deletion test/manifest.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,8 @@ import { describe, expect, it } from 'vitest'
import fs from 'fs'
import { icebergManifests } from '../src/manifest.js'
import { icebergMetadata } from '../src/metadata.js'
import { localResolver } from './helpers.js'
import { writeDataManifest } from '../src/write/manifest.js'
import { localResolver, memResolver } from './helpers.js'

describe('Iceberg Manifests', () => {
const tableUrl = 's3://hyperparam-iceberg/spark/bunnies'
Expand Down Expand Up @@ -91,4 +92,43 @@ describe('Iceberg Manifests', () => {
expect(manifests).toHaveLength(1)
expect(calls).toEqual([{ url: manifestPath, byteLength: manifestLength }])
})

it('inherits a null entry snapshot id from the manifest list', async () => {
const { resolver: memory } = memResolver()
const manifestPath = 'http://test/inherited-snapshot-id.avro'
const writer = memory.writer?.(manifestPath)
if (!writer) throw new Error('expected resolver.writer')
await writeDataManifest({
writer,
schema: { type: 'struct', 'schema-id': 0, fields: [] },
partitionSpec: { 'spec-id': 0, fields: [] },
snapshotId: /** @type {any} */ (null),
dataFiles: [{
content: 0,
file_path: 'http://test/data.parquet',
file_format: 'parquet',
partition: {},
record_count: 1n,
file_size_in_bytes: 1n,
}],
})

const metadata = /** @type {any} */ ({
'current-snapshot-id': 77,
snapshots: [{
'snapshot-id': 77,
manifests: [{
manifest_path: manifestPath,
manifest_length: BigInt(writer.offset),
partition_spec_id: 0,
content: 0,
sequence_number: 5n,
added_snapshot_id: 77n,
}],
}],
})

const manifests = await icebergManifests({ metadata, resolver: memory })
expect(manifests[0].entries[0].snapshot_id).toBe(77n)
})
})
180 changes: 178 additions & 2 deletions test/write/manifest.test.js
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
import { describe, expect, it } from 'vitest'
import { ByteWriter } from 'hyparquet-writer'
import { writeDataManifest } from '../../src/write/manifest.js'
import { writeDataManifest, writeExistingDataManifest } from '../../src/write/manifest.js'
import { avroMetadata } from '../../src/avro/avro.metadata.js'
import { avroRead } from '../../src/avro/avro.read.js'

/**
* @import {DataFile, PartitionSpec, Schema} from '../../src/types.js'
* @import {DataFile, ManifestEntry, PartitionSpec, Schema} from '../../src/types.js'
*/

describe('writeDataManifest', () => {
Expand Down Expand Up @@ -119,3 +119,179 @@ describe('writeDataManifest', () => {
expect(records[0].data_file.first_row_id).toBe(1000n)
})
})

describe('writeExistingDataManifest', () => {
/** @type {Schema} */
const schema = {
type: 'struct',
'schema-id': 0,
fields: [
{ id: 1, name: 'id', required: true, type: 'long' },
{ id: 2, name: 'name', required: false, type: 'string' },
],
}

/** @type {PartitionSpec} */
const unpartitioned = { 'spec-id': 0, fields: [] }

/** @type {DataFile} */
const dataFile = {
content: 0,
file_path: 's3://bucket/table/data/abc.parquet',
file_format: 'parquet',
partition: {},
record_count: 3n,
file_size_in_bytes: 421n,
sort_order_id: 0,
}

/** @type {ManifestEntry} */
const entry = {
status: 1,
snapshot_id: 111n,
sequence_number: 7n,
file_sequence_number: 7n,
partition_spec_id: 0,
data_file: dataFile,
}

it('writes EXISTING entries that keep their original snapshot and sequence numbers', async () => {
const writer = new ByteWriter()
writeExistingDataManifest({ writer, schema, partitionSpec: unpartitioned, entries: [entry] })
const buffer = writer.getBuffer()

const reader = { view: new DataView(buffer), offset: 0 }
const { metadata, syncMarker } = await avroMetadata(reader)
expect(metadata.content).toBe('data')

const records = await avroRead({ reader, metadata, syncMarker })
expect(records).toHaveLength(1)
// status 0 (EXISTING) with explicit numbers: an ADDED entry would leave
// these null and inherit the rewriting snapshot's sequence number, which
// would move the file forward in time and change which deletes apply.
expect(records[0]).toMatchObject({
status: 0,
snapshot_id: 111n,
sequence_number: 7n,
file_sequence_number: 7n,
data_file: { file_path: 's3://bucket/table/data/abc.parquet' },
})
})

it('carries v3 first_row_id through', async () => {
const writer = new ByteWriter()
writeExistingDataManifest({
writer,
schema,
partitionSpec: unpartitioned,
entries: [{ ...entry, data_file: { ...dataFile, first_row_id: 1000n } }],
formatVersion: 3,
})
const buffer = writer.getBuffer()

const reader = { view: new DataView(buffer), offset: 0 }
const { metadata, syncMarker } = await avroMetadata(reader)
const records = await avroRead({ reader, metadata, syncMarker })
expect(records[0].data_file.first_row_id).toBe(1000n)
})

it('round-trips stat maps from a read-decoded manifest', async () => {
// The Avro reader hands Iceberg maps back as {key, value} record arrays
// rather than plain objects, and a manifest rewrite feeds exactly those
// decoded entries back in. Sequence numbers are supplied here the way
// `icebergManifests` materializes them from the manifest list.
const first = new ByteWriter()
writeDataManifest({
writer: first,
schema,
partitionSpec: unpartitioned,
snapshotId: 111n,
dataFiles: [{
...dataFile,
value_counts: { 1: 3n, 2: 3n },
lower_bounds: { 1: new Uint8Array([1, 0, 0, 0, 0, 0, 0, 0]) },
split_offsets: [4n, 100n],
}],
})
const firstBuffer = first.getBuffer()

const firstReader = { view: new DataView(firstBuffer), offset: 0 }
const firstMeta = await avroMetadata(firstReader)
const decoded = (await avroRead({
reader: firstReader,
metadata: firstMeta.metadata,
syncMarker: firstMeta.syncMarker,
}))[0]

const second = new ByteWriter()
writeExistingDataManifest({
writer: second,
schema,
partitionSpec: unpartitioned,
entries: [/** @type {ManifestEntry} */ ({
...decoded,
sequence_number: 7n,
file_sequence_number: 7n,
partition_spec_id: 0,
})],
})
const secondBuffer = second.getBuffer()

const secondReader = { view: new DataView(secondBuffer), offset: 0 }
const secondMeta = await avroMetadata(secondReader)
const records = await avroRead({
reader: secondReader,
metadata: secondMeta.metadata,
syncMarker: secondMeta.syncMarker,
})
expect(records[0].data_file.value_counts).toEqual(decoded.data_file.value_counts)
expect(records[0].data_file.lower_bounds).toEqual(decoded.data_file.lower_bounds)
expect(records[0].data_file.split_offsets).toEqual(decoded.data_file.split_offsets)
expect(records[0].data_file.record_count).toBe(3n)
})

it('rejects deleted entries', () => {
expect(() => writeExistingDataManifest({
writer: new ByteWriter(),
schema,
partitionSpec: unpartitioned,
entries: [{ ...entry, status: 2 }],
})).toThrow('writeExistingDataManifest cannot rewrite deleted entries as existing')
})

it('rejects delete files', () => {
expect(() => writeExistingDataManifest({
writer: new ByteWriter(),
schema,
partitionSpec: unpartitioned,
entries: [{ ...entry, data_file: { ...dataFile, content: 1 } }],
})).toThrow('writeExistingDataManifest expects data files')
})

it('rejects entries without sequence numbers', () => {
expect(() => writeExistingDataManifest({
writer: new ByteWriter(),
schema,
partitionSpec: unpartitioned,
entries: [{ ...entry, file_sequence_number: undefined }],
})).toThrow('existing data manifest entry missing sequence numbers')
})

it('rejects entries without a materialized snapshot id', () => {
expect(() => writeExistingDataManifest({
writer: new ByteWriter(),
schema,
partitionSpec: unpartitioned,
entries: [{ ...entry, snapshot_id: undefined }],
})).toThrow('existing data manifest entry missing snapshot id')
})

it('rejects entries from another partition spec', () => {
expect(() => writeExistingDataManifest({
writer: new ByteWriter(),
schema,
partitionSpec: unpartitioned,
entries: [{ ...entry, partition_spec_id: 1 }],
})).toThrow('existing data entry partition spec 1 does not match 0')
})
})