From f52b0c690586b294bd6f59772f22ebcfab7e98d5 Mon Sep 17 00:00:00 2001 From: ivscheianu Date: Wed, 29 Jul 2026 07:12:29 +0300 Subject: [PATCH 1/3] feat: propagate file_format_version to CommitBuilder.storageFormat() MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit When file_format_version is set in Spark write options, lance-spark correctly encodes fragment data files in the requested format. However, at commit time the manifest's data_storage_format was never updated because CommitBuilder.storageFormat() was never called. Forward writeOptions.getFileFormatVersion() to commitBuilder.storageFormat() in all code paths that construct a CommitBuilder: - LanceBatchWrite.commit() (batch write Append/Overwrite) - StagedCommit.commitNewTable() and commitExistingTable() (staged catalog operations CREATE/REPLACE/CREATE_OR_REPLACE) - SparkPositionDeltaWrite (row-level UPDATE/DELETE/MERGE, Spark 3.4+3.5) - AddColumnsBackfillBatchWrite (add-columns backfill) - UpdateColumnsBackfillBatchWrite (column rewrite backfill) Add fileFormatVersion field to StagedCommitOptions so staged commit paths receive the value from the catalog's CreateTableSpec resolution. When fileFormatVersion is null (user didn't set the option), nothing changes — behavior is identical to before. --- .../spark/write/SparkPositionDeltaWrite.java | 4 ++ .../spark/write/SparkPositionDeltaWrite.java | 4 ++ .../spark/BaseLanceNamespaceSparkCatalog.java | 68 ++++++++++++------- .../write/AddColumnsBackfillBatchWrite.java | 16 +++-- .../lance/spark/write/LanceBatchWrite.java | 4 ++ .../org/lance/spark/write/StagedCommit.java | 8 +++ .../spark/write/StagedCommitOptions.java | 26 ++++++- .../UpdateColumnsBackfillBatchWrite.java | 16 +++-- .../spark/write/LanceBatchWriteTest.java | 4 +- .../spark/write/StagedCommitOptionsTest.java | 11 ++- .../lance/spark/write/StagedCommitTest.java | 26 ++++--- 11 files changed, 134 insertions(+), 53 deletions(-) diff --git a/lance-spark-3.4_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java b/lance-spark-3.4_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java index 6e6344de0..576b24cd2 100644 --- a/lance-spark-3.4_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java +++ b/lance-spark-3.4_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java @@ -207,6 +207,10 @@ public void commit(WriterCommitMessage[] messages) { .writeParams( LanceRuntime.mergeStorageOptions( writeOptions.getStorageOptions(), initialStorageOptions)); + String fileFormatVersion = writeOptions.getFileFormatVersion(); + if (fileFormatVersion != null) { + commitBuilder.storageFormat(fileFormatVersion); + } if (dataset.hasStableRowIds() || Boolean.TRUE.equals(writeOptions.getEnableStableRowIds())) { commitBuilder.useStableRowIds(true); diff --git a/lance-spark-3.5_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java b/lance-spark-3.5_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java index e1b867836..707a1642b 100644 --- a/lance-spark-3.5_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java +++ b/lance-spark-3.5_2.12/src/main/java/org/lance/spark/write/SparkPositionDeltaWrite.java @@ -228,6 +228,10 @@ public void commit(WriterCommitMessage[] messages) { .writeParams( LanceRuntime.mergeStorageOptions( writeOptions.getStorageOptions(), initialStorageOptions)); + String fileFormatVersion = writeOptions.getFileFormatVersion(); + if (fileFormatVersion != null) { + commitBuilder.storageFormat(fileFormatVersion); + } if (useStableRowIds) { commitBuilder.useStableRowIds(true); } diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/BaseLanceNamespaceSparkCatalog.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/BaseLanceNamespaceSparkCatalog.java index 435e1b6eb..91bdf99b7 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/BaseLanceNamespaceSparkCatalog.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/BaseLanceNamespaceSparkCatalog.java @@ -1125,6 +1125,7 @@ public StagedTable stageCreate( StagedCommitOptions.of( merged, catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion, namespace, tableIdList, managedVersioning); @@ -1161,7 +1162,9 @@ private StagedTable stageCreateAtPath( Schema arrowSchema = LanceArrowUtils.toArrowSchema(processedSchema, "UTC", true); final StagedCommitOptions commitOptions = StagedCommitOptions.pathBased( - catalogConfig.getStorageOptions(), catalogConfig.isEnableStableRowIds(properties)); + catalogConfig.getStorageOptions(), + catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion); StagedCommit stagedCommit = StagedCommit.forNewTable(arrowSchema, datasetUri, commitOptions); stagedCommit.setShardingSpec(shardingSpec); return createStagedDataset( @@ -1201,19 +1204,20 @@ public StagedTable stageReplace( Dataset ds = Utils.openDatasetBuilder(resolved.readOptions).build(); Map merged = LanceRuntime.mergeStorageOptions(catalogConfig.getStorageOptions(), initialStorageOptions); + // Use specified file format version, or fall back to existing table's version + if (fileFormatVersion == null) { + fileFormatVersion = ds.getLanceFileFormatVersion(); + } final StagedCommitOptions commitOptions = StagedCommitOptions.of( merged, catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion, namespace, resolved.tableIdList, managedVersioning); StagedCommit stagedCommit = StagedCommit.forExistingTable(ds, arrowSchema, commitOptions); stagedCommit.setShardingSpec(shardingSpec); - // Use specified file format version, or fall back to existing table's version - if (fileFormatVersion == null) { - fileFormatVersion = ds.getLanceFileFormatVersion(); - } return createStagedDataset( resolved.readOptions, processedSchema, @@ -1250,16 +1254,18 @@ private StagedTable stageReplaceAtPath( throw new NoSuchTableException(ident); } + // Use specified file format version, or fall back to existing table's version + if (fileFormatVersion == null) { + fileFormatVersion = ds.getLanceFileFormatVersion(); + } Schema arrowSchema = LanceArrowUtils.toArrowSchema(processedSchema, "UTC", true); final StagedCommitOptions commitOptions = StagedCommitOptions.pathBased( - catalogConfig.getStorageOptions(), catalogConfig.isEnableStableRowIds(properties)); + catalogConfig.getStorageOptions(), + catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion); StagedCommit stagedCommit = StagedCommit.forExistingTable(ds, arrowSchema, commitOptions); stagedCommit.setShardingSpec(shardingSpec); - // Use specified file format version, or fall back to existing table's version - if (fileFormatVersion == null) { - fileFormatVersion = ds.getLanceFileFormatVersion(); - } return createStagedDataset( readOptions, processedSchema, @@ -1333,24 +1339,32 @@ public StagedTable stageCreateOrReplace( name); Schema arrowSchema = LanceArrowUtils.toArrowSchema(processedSchema, "UTC", true); - // Use specified file format version, or fall back to existing table's version Map merged = LanceRuntime.mergeStorageOptions(catalogConfig.getStorageOptions(), initialStorageOptions); - final StagedCommitOptions commitOptions = - StagedCommitOptions.of( - merged, - catalogConfig.isEnableStableRowIds(properties), - namespace, - tableIdList, - managedVersioning); StagedCommit stagedCommit; if (exists) { Dataset ds = Utils.openDatasetBuilder(readOptions).build(); - stagedCommit = StagedCommit.forExistingTable(ds, arrowSchema, commitOptions); if (fileFormatVersion == null) { fileFormatVersion = ds.getLanceFileFormatVersion(); } + final StagedCommitOptions commitOptions = + StagedCommitOptions.of( + merged, + catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion, + namespace, + tableIdList, + managedVersioning); + stagedCommit = StagedCommit.forExistingTable(ds, arrowSchema, commitOptions); } else { + final StagedCommitOptions commitOptions = + StagedCommitOptions.of( + merged, + catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion, + namespace, + tableIdList, + managedVersioning); stagedCommit = StagedCommit.forNewTable(arrowSchema, location, commitOptions); } stagedCommit.setShardingSpec(shardingSpec); @@ -1384,18 +1398,24 @@ private StagedTable stageCreateOrReplaceAtPath( boolean exists = tableExistsAtPath(ident); Schema arrowSchema = LanceArrowUtils.toArrowSchema(processedSchema, "UTC", true); - final StagedCommitOptions commitOptions = - StagedCommitOptions.pathBased( - catalogConfig.getStorageOptions(), catalogConfig.isEnableStableRowIds(properties)); StagedCommit stagedCommit; - // Use specified file format version, or fall back to existing table's version if (exists) { Dataset ds = Utils.openDatasetBuilder(readOptions).build(); - stagedCommit = StagedCommit.forExistingTable(ds, arrowSchema, commitOptions); if (fileFormatVersion == null) { fileFormatVersion = ds.getLanceFileFormatVersion(); } + final StagedCommitOptions commitOptions = + StagedCommitOptions.pathBased( + catalogConfig.getStorageOptions(), + catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion); + stagedCommit = StagedCommit.forExistingTable(ds, arrowSchema, commitOptions); } else { + final StagedCommitOptions commitOptions = + StagedCommitOptions.pathBased( + catalogConfig.getStorageOptions(), + catalogConfig.isEnableStableRowIds(properties), + fileFormatVersion); stagedCommit = StagedCommit.forNewTable(arrowSchema, datasetUri, commitOptions); } stagedCommit.setShardingSpec(shardingSpec); diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/AddColumnsBackfillBatchWrite.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/AddColumnsBackfillBatchWrite.java index 7615e6218..d59050c18 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/AddColumnsBackfillBatchWrite.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/AddColumnsBackfillBatchWrite.java @@ -149,14 +149,18 @@ public void commit(WriterCommitMessage[] messages) { // Commit merge operation using CommitBuilder Merge merge = Merge.builder().fragments(fragments).schema(arrowSchema).build(); + CommitBuilder commitBuilder = + new CommitBuilder(dataset) + .writeParams( + LanceRuntime.mergeStorageOptions( + writeOptions.getStorageOptions(), initialStorageOptions)); + String fileFormatVersion = writeOptions.getFileFormatVersion(); + if (fileFormatVersion != null) { + commitBuilder.storageFormat(fileFormatVersion); + } try (Transaction txn = new Transaction.Builder().readVersion(version).operation(merge).build(); - Dataset committed = - new CommitBuilder(dataset) - .writeParams( - LanceRuntime.mergeStorageOptions( - writeOptions.getStorageOptions(), initialStorageOptions)) - .execute(txn)) { + Dataset committed = commitBuilder.execute(txn)) { // auto-close txn and committed dataset } } diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java index 5cba05f2e..e33657071 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java @@ -202,6 +202,10 @@ public void commit(WriterCommitMessage[] messages) { .writeParams( LanceRuntime.mergeStorageOptions( writeOptions.getStorageOptions(), initialStorageOptions)); + String fileFormatVersion = writeOptions.getFileFormatVersion(); + if (fileFormatVersion != null) { + commitBuilder.storageFormat(fileFormatVersion); + } // When enableStableRowIds is null (user didn't pass the option), // lance-core auto-inherits the flag from the existing manifest. // Appending to a table with stable row IDs works without diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java index dd6423b70..cc84be599 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java @@ -53,6 +53,7 @@ public class StagedCommit { // The non-staged path (LanceBatchWrite) uses boxed Boolean because null means // "user didn't specify" and lets lance-core inherit the flag from the manifest. private boolean enableStableRowIds; + private final String fileFormatVersion; private List fragments; private Schema schema; private ShardingSpec shardingSpec; @@ -95,6 +96,7 @@ private StagedCommit( this.storageOptions = new HashMap<>(options.getStorageOptions()); this.isNewTable = datasetUri != null; this.enableStableRowIds = options.isEnableStableRowIds(); + this.fileFormatVersion = options.getFileFormatVersion(); this.namespace = options.getNamespace(); this.tableId = options.getTableId(); this.managedVersioning = options.isManagedVersioning(); @@ -150,6 +152,9 @@ private void commitNewTable() { if (enableStableRowIds) { builder.useStableRowIds(true); } + if (fileFormatVersion != null) { + builder.storageFormat(fileFormatVersion); + } applyManagedVersioning(builder); try (Transaction txn = new Transaction.Builder().operation(operation).build(); Dataset committed = builder.execute(txn)) { @@ -167,6 +172,9 @@ private void commitExistingTable() { final CommitBuilder builder = new CommitBuilder(uri, LanceRuntime.allocator()).writeParams(storageOptions); builder.useStableRowIds(enableStableRowIds); + if (fileFormatVersion != null) { + builder.storageFormat(fileFormatVersion); + } applyManagedVersioning(builder); try (Dataset committed = commitOperation(builder, version, operation)) { SparkLanceShardingUtils.initializeMemWal(committed, shardingSpec); diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommitOptions.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommitOptions.java index 8807b320a..229fbfd0f 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommitOptions.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommitOptions.java @@ -27,6 +27,7 @@ public class StagedCommitOptions { private final Map storageOptions; private final boolean enableStableRowIds; + private final String fileFormatVersion; private final LanceNamespace namespace; private final List tableId; private final boolean managedVersioning; @@ -34,11 +35,13 @@ public class StagedCommitOptions { private StagedCommitOptions( final Map storageOptions, final boolean enableStableRowIds, + final String fileFormatVersion, final LanceNamespace namespace, final List tableId, final boolean managedVersioning) { this.storageOptions = storageOptions; this.enableStableRowIds = enableStableRowIds; + this.fileFormatVersion = fileFormatVersion; this.namespace = namespace; this.tableId = tableId; this.managedVersioning = managedVersioning; @@ -47,17 +50,30 @@ private StagedCommitOptions( public static StagedCommitOptions of( final Map storageOptions, final boolean enableStableRowIds, + final String fileFormatVersion, final LanceNamespace namespace, final List tableId, final boolean managedVersioning) { return new StagedCommitOptions( - storageOptions, enableStableRowIds, namespace, tableId, managedVersioning); + storageOptions, + enableStableRowIds, + fileFormatVersion, + namespace, + tableId, + managedVersioning); } public static StagedCommitOptions pathBased( - final Map storageOptions, final boolean enableStableRowIds) { + final Map storageOptions, + final boolean enableStableRowIds, + final String fileFormatVersion) { return new StagedCommitOptions( - storageOptions, enableStableRowIds, NO_NAMESPACE, NO_TABLE_ID, NO_MANAGED_VERSIONING); + storageOptions, + enableStableRowIds, + fileFormatVersion, + NO_NAMESPACE, + NO_TABLE_ID, + NO_MANAGED_VERSIONING); } public Map getStorageOptions() { @@ -68,6 +84,10 @@ public boolean isEnableStableRowIds() { return enableStableRowIds; } + public String getFileFormatVersion() { + return fileFormatVersion; + } + public LanceNamespace getNamespace() { return namespace; } diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/UpdateColumnsBackfillBatchWrite.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/UpdateColumnsBackfillBatchWrite.java index ac4271815..967a9c58b 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/UpdateColumnsBackfillBatchWrite.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/UpdateColumnsBackfillBatchWrite.java @@ -156,14 +156,18 @@ public void commit(WriterCommitMessage[] messages) { Objects.requireNonNull( writeOptions.getVersion(), "version must be set (resolved in UpdateColumnsBackfillBatchWrite constructor)"); + CommitBuilder commitBuilder = + new CommitBuilder(dataset) + .writeParams( + LanceRuntime.mergeStorageOptions( + writeOptions.getStorageOptions(), initialStorageOptions)); + String fileFormatVersion = writeOptions.getFileFormatVersion(); + if (fileFormatVersion != null) { + commitBuilder.storageFormat(fileFormatVersion); + } try (Transaction txn = new Transaction.Builder().readVersion(version).operation(update).build(); - Dataset committed = - new CommitBuilder(dataset) - .writeParams( - LanceRuntime.mergeStorageOptions( - writeOptions.getStorageOptions(), initialStorageOptions)) - .execute(txn)) { + Dataset committed = commitBuilder.execute(txn)) { // auto-close txn and committed dataset } } diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java index d60daf7f2..80e8f12ac 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java @@ -155,7 +155,7 @@ public void testCommitMergesInitialStorageOptionsIntoStagedCommit(TestInfo testI // StagedCommit starts with no storage options, as for a path-based staged create. StagedCommit stagedCommit = StagedCommit.forNewTable( - schema, datasetUri, StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + schema, datasetUri, StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); // initialStorageOptions represents namespace-vended credentials, passed to LanceBatchWrite // independently of how stagedCommit was constructed. @@ -207,7 +207,7 @@ public void testCommitStagedMergePrefersInitialStorageOptionsOnConflict(TestInfo StagedCommit stagedCommit = StagedCommit.forNewTable( - schema, datasetUri, StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + schema, datasetUri, StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); LanceBatchWrite batchWrite = new LanceBatchWrite( diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitOptionsTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitOptionsTest.java index 908e7fc4c..b8955a4a6 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitOptionsTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitOptionsTest.java @@ -49,12 +49,14 @@ public void testOfRetainsAllFields() { StagedCommitOptions.of( storageOptions, STABLE_ROW_IDS_ENABLED, + "2.2", NO_NAMESPACE, tableId, MANAGED_VERSIONING_ENABLED); assertEquals(storageOptions, options.getStorageOptions()); assertTrue(options.isEnableStableRowIds()); + assertEquals("2.2", options.getFileFormatVersion()); assertNull(options.getNamespace()); assertEquals(tableId, options.getTableId()); assertTrue(options.isManagedVersioning()); @@ -69,6 +71,7 @@ public void testOfWithStableRowIdsDisabled() { StagedCommitOptions.of( storageOptions, STABLE_ROW_IDS_DISABLED, + null, NO_NAMESPACE, tableId, MANAGED_VERSIONING_DISABLED); @@ -82,10 +85,11 @@ public void testPathBasedSetsDefaults() { final Map storageOptions = Collections.singletonMap("path", "/tmp/lance"); final StagedCommitOptions options = - StagedCommitOptions.pathBased(storageOptions, STABLE_ROW_IDS_ENABLED); + StagedCommitOptions.pathBased(storageOptions, STABLE_ROW_IDS_ENABLED, "2.1"); assertEquals(storageOptions, options.getStorageOptions()); assertTrue(options.isEnableStableRowIds()); + assertEquals("2.1", options.getFileFormatVersion()); assertNull(options.getNamespace()); assertNull(options.getTableId()); assertFalse(options.isManagedVersioning()); @@ -96,7 +100,7 @@ public void testPathBasedWithStableRowIdsDisabled() { final Map storageOptions = Collections.emptyMap(); final StagedCommitOptions options = - StagedCommitOptions.pathBased(storageOptions, STABLE_ROW_IDS_DISABLED); + StagedCommitOptions.pathBased(storageOptions, STABLE_ROW_IDS_DISABLED, null); assertFalse(options.isEnableStableRowIds()); assertNull(options.getNamespace()); @@ -110,6 +114,7 @@ public void testOfWithEmptyStorageOptions() { StagedCommitOptions.of( Collections.emptyMap(), STABLE_ROW_IDS_DISABLED, + null, NO_NAMESPACE, NO_TABLE_ID, MANAGED_VERSIONING_DISABLED); @@ -124,7 +129,7 @@ public void testStorageOptionsAreNotCopied() { storageOptions.put("key", "original"); final StagedCommitOptions options = - StagedCommitOptions.pathBased(storageOptions, STABLE_ROW_IDS_DISABLED); + StagedCommitOptions.pathBased(storageOptions, STABLE_ROW_IDS_DISABLED, null); storageOptions.put("key", "modified"); assertEquals("modified", options.getStorageOptions().get("key")); diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java index 040a0c33c..4fdcf8794 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java @@ -60,7 +60,9 @@ public void testCommitNewTable(TestInfo testInfo) { TestUtils.getDatasetUri(tempDir.toString(), testInfo.getTestMethod().get().getName()); StagedCommit commit = StagedCommit.forNewTable( - ARROW_SCHEMA, datasetUri, StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + ARROW_SCHEMA, + datasetUri, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); commit.commit(); try (Dataset dataset = Dataset.open(datasetUri, LanceRuntime.allocator())) { assertEquals(0, dataset.countRows()); @@ -73,7 +75,9 @@ public void testCommitExistingTable(TestInfo testInfo) { Dataset dataset = Dataset.open(datasetUri, LanceRuntime.allocator()); StagedCommit commit = StagedCommit.forExistingTable( - dataset, ARROW_SCHEMA, StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + dataset, + ARROW_SCHEMA, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); commit.commit(); try (Dataset reopened = Dataset.open(datasetUri, LanceRuntime.allocator())) { assertEquals(0, reopened.countRows()); @@ -86,7 +90,9 @@ public void testAbortNewTableWithoutNamespace(TestInfo testInfo) { TestUtils.getDatasetUri(tempDir.toString(), testInfo.getTestMethod().get().getName()); StagedCommit commit = StagedCommit.forNewTable( - ARROW_SCHEMA, datasetUri, StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + ARROW_SCHEMA, + datasetUri, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); commit.abort(); } @@ -96,7 +102,9 @@ public void testAbortExistingTableClosesDataset(TestInfo testInfo) { Dataset dataset = Dataset.open(datasetUri, LanceRuntime.allocator()); StagedCommit commit = StagedCommit.forExistingTable( - dataset, ARROW_SCHEMA, StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + dataset, + ARROW_SCHEMA, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); commit.abort(); } @@ -106,7 +114,7 @@ public void testMergeStorageOptionsAddsNewKeys() { StagedCommit.forNewTable( ARROW_SCHEMA, "unused://uri", - StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); Map extra = new HashMap<>(); extra.put("access_key_id", "AKIA..."); @@ -123,7 +131,7 @@ public void testMergeStorageOptionsOverridesExistingKeys() { base.put("region", "us-west-2"); StagedCommit commit = StagedCommit.forNewTable( - ARROW_SCHEMA, "unused://uri", StagedCommitOptions.pathBased(base, false)); + ARROW_SCHEMA, "unused://uri", StagedCommitOptions.pathBased(base, false, null)); commit.mergeStorageOptions(Collections.singletonMap("access_key_id", "fresh-key")); @@ -136,7 +144,7 @@ public void testMergeStorageOptionsNullAndEmptyAreNoOps() { Map base = Collections.singletonMap("access_key_id", "AKIA..."); StagedCommit commit = StagedCommit.forNewTable( - ARROW_SCHEMA, "unused://uri", StagedCommitOptions.pathBased(base, false)); + ARROW_SCHEMA, "unused://uri", StagedCommitOptions.pathBased(base, false, null)); commit.mergeStorageOptions(null); commit.mergeStorageOptions(Collections.emptyMap()); @@ -150,7 +158,7 @@ public void testMergeStorageOptionsDoesNotMutateCallerMap() { StagedCommit.forNewTable( ARROW_SCHEMA, "unused://uri", - StagedCommitOptions.pathBased(Collections.emptyMap(), false)); + StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); Map extra = new HashMap<>(); extra.put("access_key_id", "AKIA..."); @@ -171,7 +179,7 @@ public void testMergeStorageOptionsAcceptsUnmodifiableBaseMap() { StagedCommit.forNewTable( ARROW_SCHEMA, "unused://uri", - StagedCommitOptions.pathBased(unmodifiableCatalogOptions, false)); + StagedCommitOptions.pathBased(unmodifiableCatalogOptions, false, null)); assertDoesNotThrow( () -> commit.mergeStorageOptions(Collections.singletonMap("access_key_id", "AKIA..."))); From fd3900d67ea2c3f32624b0b05e51151460b48778 Mon Sep 17 00:00:00 2001 From: ivscheianu Date: Thu, 6 Aug 2026 09:42:00 +0300 Subject: [PATCH 2/3] chore: bump lance-core to 11.0.0-beta.1 --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 363f90a5f..e32b44835 100644 --- a/pom.xml +++ b/pom.xml @@ -51,7 +51,7 @@ 0.6.1 - 9.0.0 + 11.0.0-beta.1 0.8.6 0.4.0 From 5ff1c3c01f8693d0889024e37a7072f5eb51bc94 Mon Sep 17 00:00:00 2001 From: ivscheianu Date: Mon, 10 Aug 2026 16:31:02 +0300 Subject: [PATCH 3/3] fix: propagate write-time file_format_version to StagedCommit When DataFrameWriterV2.option("file_format_version", ...) sets the format at write time, LanceBatchWrite.commit() now forwards it to StagedCommit via the new setFileFormatVersion setter. Write-time takes precedence over the stage-time value; null write-time preserves whatever the catalog resolved at stage time. Without this, fragments are encoded in the requested format but the manifest's data_storage_format remains at the default, causing check_storage_version to reject the commit. --- .../lance/spark/write/LanceBatchWrite.java | 4 + .../org/lance/spark/write/StagedCommit.java | 11 +- .../spark/write/LanceBatchWriteTest.java | 101 ++++++++++++++++++ .../lance/spark/write/StagedCommitTest.java | 31 ++++++ 4 files changed, 146 insertions(+), 1 deletion(-) diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java index e33657071..d3792dac8 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/LanceBatchWrite.java @@ -178,6 +178,10 @@ public void commit(WriterCommitMessage[] messages) { if (enableStableRowIds != null) { stagedCommit.setEnableStableRowIds(enableStableRowIds); } + String fileFormatVersion = writeOptions.getFileFormatVersion(); + if (fileFormatVersion != null) { + stagedCommit.setFileFormatVersion(fileFormatVersion); + } // For a path-based staged create, StagedCommit only has the catalog's static storage // options at this point. Merge in write-time and namespace-vended options now so // StagedCommit.commit() uses them. Mirrors the non-staged merge below. diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java index cc84be599..453f3e956 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/write/StagedCommit.java @@ -53,7 +53,7 @@ public class StagedCommit { // The non-staged path (LanceBatchWrite) uses boxed Boolean because null means // "user didn't specify" and lets lance-core inherit the flag from the manifest. private boolean enableStableRowIds; - private final String fileFormatVersion; + private String fileFormatVersion; private List fragments; private Schema schema; private ShardingSpec shardingSpec; @@ -114,6 +114,10 @@ public void setEnableStableRowIds(boolean enableStableRowIds) { this.enableStableRowIds = enableStableRowIds; } + public void setFileFormatVersion(String fileFormatVersion) { + this.fileFormatVersion = fileFormatVersion; + } + public void setShardingSpec(ShardingSpec shardingSpec) { this.shardingSpec = shardingSpec; } @@ -136,6 +140,11 @@ Map getStorageOptions() { return storageOptions; } + /** Returns the current file format version. Visible for testing. */ + String getFileFormatVersion() { + return fileFormatVersion; + } + /** Performs the actual commit using the stored dataset and fragments. */ public void commit() { if (dataset.isEmpty()) { diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java index 80e8f12ac..130dc0b8b 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/LanceBatchWriteTest.java @@ -45,6 +45,7 @@ import java.util.Map; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; public class LanceBatchWriteTest { @@ -229,6 +230,106 @@ public void testCommitStagedMergePrefersInitialStorageOptionsOnConflict(TestInfo assertEquals("us-west-2", stagedCommit.getStorageOptions().get("region")); } + /** + * When a write-time option sets {@code file_format_version} (e.g. via {@code + * DataFrameWriterV2.option("file_format_version", "2.2")}), the staged commit must use that + * version even if none was set at stage time. + */ + @Test + public void testCommitStagedPropagatesWriteTimeFileFormatVersion(TestInfo testInfo) { + String datasetName = testInfo.getTestMethod().get().getName(); + String datasetUri = TestUtils.getDatasetUri(tempDir.toString(), datasetName); + + Field field = new Field("column1", FieldType.nullable(new ArrowType.Int(32, true)), null); + Schema schema = new Schema(Collections.singletonList(field)); + StructType sparkSchema = LanceArrowUtils.fromArrowSchema(schema); + + // Stage-time: no file format version set (simulates stageCreate without table property) + StagedCommit stagedCommit = + StagedCommit.forNewTable( + schema, datasetUri, StagedCommitOptions.pathBased(Collections.emptyMap(), false, null)); + assertNull(stagedCommit.getFileFormatVersion()); + + // Write-time: user set file_format_version via DataFrameWriterV2.option(...) + LanceSparkWriteOptions writeOptions = + LanceSparkWriteOptions.from(datasetUri).toBuilder().fileFormatVersion("2.1").build(); + + LanceBatchWrite batchWrite = + new LanceBatchWrite( + sparkSchema, writeOptions, false, null, null, null, null, false, stagedCommit); + + batchWrite.commit(new WriterCommitMessage[0]); + + assertEquals("2.1", stagedCommit.getFileFormatVersion()); + } + + /** + * Write-time {@code file_format_version} takes precedence over the stage-time value. This covers + * the case where the catalog resolved a default at stage time but the user explicitly overrides + * it at write time. + */ + @Test + public void testCommitStagedWriteTimeFileFormatVersionTakesPrecedence(TestInfo testInfo) { + String datasetName = testInfo.getTestMethod().get().getName(); + String datasetUri = TestUtils.getDatasetUri(tempDir.toString(), datasetName); + + Field field = new Field("column1", FieldType.nullable(new ArrowType.Int(32, true)), null); + Schema schema = new Schema(Collections.singletonList(field)); + StructType sparkSchema = LanceArrowUtils.fromArrowSchema(schema); + + // Stage-time: catalog resolved file format version "2.0" + StagedCommit stagedCommit = + StagedCommit.forNewTable( + schema, + datasetUri, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, "2.0")); + assertEquals("2.0", stagedCommit.getFileFormatVersion()); + + // Write-time: user explicitly overrides to "2.2" + LanceSparkWriteOptions writeOptions = + LanceSparkWriteOptions.from(datasetUri).toBuilder().fileFormatVersion("2.2").build(); + + LanceBatchWrite batchWrite = + new LanceBatchWrite( + sparkSchema, writeOptions, false, null, null, null, null, false, stagedCommit); + + batchWrite.commit(new WriterCommitMessage[0]); + + assertEquals("2.2", stagedCommit.getFileFormatVersion()); + } + + /** + * When write options have no {@code file_format_version}, the stage-time value must be preserved. + * Null write-time must not overwrite a stage-time resolution. + */ + @Test + public void testCommitStagedPreservesStageTimeVersionWhenWriteTimeIsNull(TestInfo testInfo) { + String datasetName = testInfo.getTestMethod().get().getName(); + String datasetUri = TestUtils.getDatasetUri(tempDir.toString(), datasetName); + + Field field = new Field("column1", FieldType.nullable(new ArrowType.Int(32, true)), null); + Schema schema = new Schema(Collections.singletonList(field)); + StructType sparkSchema = LanceArrowUtils.fromArrowSchema(schema); + + // Stage-time: version set from table properties + StagedCommit stagedCommit = + StagedCommit.forNewTable( + schema, + datasetUri, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, "2.1")); + + // Write-time: no file_format_version option + LanceSparkWriteOptions writeOptions = LanceSparkWriteOptions.from(datasetUri); + + LanceBatchWrite batchWrite = + new LanceBatchWrite( + sparkSchema, writeOptions, false, null, null, null, null, false, stagedCommit); + + batchWrite.commit(new WriterCommitMessage[0]); + + assertEquals("2.1", stagedCommit.getFileFormatVersion()); + } + private static WriterCommitMessage writeRows( LanceBatchWrite batchWrite, StructType sparkSchema, int numRows, int startValue) throws Exception { diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java index 4fdcf8794..a25384a15 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/write/StagedCommitTest.java @@ -168,6 +168,37 @@ public void testMergeStorageOptionsDoesNotMutateCallerMap() { assertEquals("AKIA...", extra.get("access_key_id")); } + @Test + public void testCommitNewTableWithFileFormatVersion(TestInfo testInfo) { + String datasetUri = + TestUtils.getDatasetUri(tempDir.toString(), testInfo.getTestMethod().get().getName()); + StagedCommit commit = + StagedCommit.forNewTable( + ARROW_SCHEMA, + datasetUri, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, "2.1")); + commit.commit(); + try (Dataset dataset = Dataset.open(datasetUri, LanceRuntime.allocator())) { + assertEquals("2.1", dataset.getLanceFileFormatVersion()); + } + } + + @Test + public void testSetFileFormatVersionOverridesStageTime(TestInfo testInfo) { + String datasetUri = + TestUtils.getDatasetUri(tempDir.toString(), testInfo.getTestMethod().get().getName()); + StagedCommit commit = + StagedCommit.forNewTable( + ARROW_SCHEMA, + datasetUri, + StagedCommitOptions.pathBased(Collections.emptyMap(), false, "2.0")); + commit.setFileFormatVersion("2.1"); + commit.commit(); + try (Dataset dataset = Dataset.open(datasetUri, LanceRuntime.allocator())) { + assertEquals("2.1", dataset.getLanceFileFormatVersion()); + } + } + @Test public void testMergeStorageOptionsAcceptsUnmodifiableBaseMap() { // Path-based staged creates pass catalogConfig.getStorageOptions() directly, which