diff --git a/src/manifest.js b/src/manifest.js index ba5a83e..629db76 100644 --- a/src/manifest.js +++ b/src/manifest.js @@ -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, diff --git a/src/write/manifest.js b/src/write/manifest.js index f27307c..047b264 100644 --- a/src/write/manifest.js +++ b/src/write/manifest.js @@ -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) { @@ -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} 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 @@ -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) { diff --git a/test/manifest.test.js b/test/manifest.test.js index a381021..744c64a 100644 --- a/test/manifest.test.js +++ b/test/manifest.test.js @@ -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' @@ -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) + }) }) diff --git a/test/write/manifest.test.js b/test/write/manifest.test.js index fca6ee7..70a75b3 100644 --- a/test/write/manifest.test.js +++ b/test/write/manifest.test.js @@ -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', () => { @@ -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') + }) +})