diff --git a/docs/src/config.md b/docs/src/config.md index 84dcfd75b..8fbb56e50 100644 --- a/docs/src/config.md +++ b/docs/src/config.md @@ -56,6 +56,7 @@ The following features require the Lance Spark SQL extension to be enabled: - [ADD COLUMNS with backfill](operations/dml/add-columns.md) - Add new columns and backfill existing rows with data - [UPDATE COLUMNS with backfill](operations/dml/update-columns.md) - Update existing columns using data from a source - [OPTIMIZE](operations/ddl/optimize.md) - Compact table fragments for improved query performance +- [OPTIMIZE INDEX](operations/ddl/optimize-index.md) - Incrementally maintain a named index - [VACUUM](operations/ddl/vacuum.md) - Remove old versions and reclaim storage space ## Basic Setup diff --git a/docs/src/operations/ddl/.pages b/docs/src/operations/ddl/.pages index bd8c6daf1..c4227cb16 100644 --- a/docs/src/operations/ddl/.pages +++ b/docs/src/operations/ddl/.pages @@ -12,6 +12,7 @@ nav: - drop-table.md - create-index.md - show-indexes.md + - optimize-index.md - create-branch.md - drop-branch.md - show-branches.md diff --git a/docs/src/operations/ddl/create-index.md b/docs/src/operations/ddl/create-index.md index 327cfeb2e..c3d9370be 100755 --- a/docs/src/operations/ddl/create-index.md +++ b/docs/src/operations/ddl/create-index.md @@ -266,12 +266,12 @@ to scanning the data until it is populated. There are two ways to populate it: ALTER TABLE lance.db.users CREATE INDEX idx_id USING zonemap (id); ``` -- **Incremental build through the SDK:** when only some fragments are unindexed (for example after - appending data to an already-built index), `Dataset.optimizeIndices` indexes just the unindexed - fragments. This currently runs on a single node: +- **Incremental maintenance:** when only some fragments are unindexed (for example after appending + data to an already-built index), [`OPTIMIZE INDEX`](optimize-index.md) indexes the uncovered + fragments. This currently runs on the Spark driver: - ```java - dataset.optimizeIndices(OptimizeOptions.builder().build()); + ```sql + ALTER TABLE lance.db.users OPTIMIZE INDEX idx_id; ``` `train = false` is supported for all index methods. Because deferred index creation does not build @@ -312,4 +312,4 @@ The `CREATE INDEX` command operates as follows: - **Index Methods**: The `zonemap`, `bitmap`, `label_list`, `ngram`, `bloomfilter`, `rtree`, `btree`, and `fts` (or `inverted`) methods are supported for index creation. - **Indexed Column Count**: All supported index methods currently support exactly one indexed column. - **Index Replacement**: If you create an index with the same name as an existing one, the old index will be replaced by the new one. -- **Deferred Training**: With `train = false` the index is registered empty and is populated later, either by re-running `CREATE INDEX` (a full distributed build that replaces the empty index) or, for incremental coverage of newly appended fragments, by `Dataset.optimizeIndices` in the SDK. The SQL `OPTIMIZE` command compacts fragments and does not train deferred indexes. +- **Deferred Training**: With `train = false` the index is registered empty and is populated later, either by re-running `CREATE INDEX` (a full distributed build that replaces the empty index) or through `ALTER TABLE ... OPTIMIZE INDEX` for driver-side incremental maintenance. The table-level `OPTIMIZE` command compacts fragments and does not train deferred indexes. diff --git a/docs/src/operations/ddl/optimize-index.md b/docs/src/operations/ddl/optimize-index.md new file mode 100644 index 000000000..6394011d4 --- /dev/null +++ b/docs/src/operations/ddl/optimize-index.md @@ -0,0 +1,61 @@ +# OPTIMIZE INDEX + +Incrementally maintains an existing named Lance index. + +!!! warning "Spark Extension Required" + This feature requires the Lance Spark SQL extension to be enabled. See [Spark SQL Extensions](../../config.md#spark-sql-extensions) for configuration details. + +## Syntax + +```sql +ALTER TABLE table_name OPTIMIZE INDEX index_name +[WITH ( + num_indices_to_merge = non_negative_integer +)]; +``` + +The index must already exist. Lance builds index data for fragments not currently covered by the +named index and may merge existing index segments according to the supplied options. + +## Options + +| Option | Type | Description | +|--------|------|-------------| +| `num_indices_to_merge` | Integer | Number of existing index segments Lance should merge during maintenance. When omitted, Lance chooses its default. Set to `0` to add coverage without requesting a merge of existing segments. | + +The option is passed directly to Lance's `OptimizeOptions`. The target index name is passed as +the sole entry in `indexNames`, so other indexes on the table are not maintained by this command. + +## Examples + +Maintain a scalar index after new data is appended: + +```sql +ALTER TABLE lance.db.users OPTIMIZE INDEX idx_user_id; +``` + +Build coverage for new fragments without requesting a merge of existing segments: + +```sql +ALTER TABLE lance.db.users OPTIMIZE INDEX idx_user_id WITH ( + num_indices_to_merge = 0 +); +``` + +## Output + +| Column | Type | Description | +|--------|------|-------------| +| `index_name` | String | Name of the maintained index. | +| `fragments_indexed` | Long | Number of previously uncovered live fragments indexed by the operation. | +| `segments_before` | Long | Number of physical segments for the named index before maintenance. | +| `segments_after` | Long | Number of physical segments for the named index after maintenance. | + +## Execution + +This command currently invokes Lance index maintenance on the Spark driver. Its SQL contract is +independent of execution strategy, so a future distributed implementation can retain the same +syntax and options. + +`ALTER TABLE ... OPTIMIZE INDEX` maintains index coverage. The table-level [`OPTIMIZE`](optimize.md) +command compacts data fragments and is a separate operation. diff --git a/integration-tests/test_lance_spark.py b/integration-tests/test_lance_spark.py index 9e8fa3f6d..d20e555be 100644 --- a/integration-tests/test_lance_spark.py +++ b/integration-tests/test_lance_spark.py @@ -915,6 +915,39 @@ def test_create_distributed_bitmap_index(self, spark): ) assert spark.sql("SELECT * FROM default.test_table").count() == 4 + def test_optimize_index(self, spark): + """Test incremental index maintenance through Spark SQL.""" + spark.sql("CREATE TABLE default.test_table (id INT, name STRING)") + spark.sql("INSERT INTO default.test_table VALUES (1, 'one'), (2, 'two')") + spark.sql(""" + ALTER TABLE default.test_table + CREATE INDEX idx_id USING zonemap (id) + """) + spark.sql("INSERT INTO default.test_table VALUES (3, 'three')") + + before = next( + row + for row in spark.sql("SHOW INDEXES IN default.test_table").collect() + if row.name == "idx_id" + ) + assert before.num_unindexed_fragments > 0 + + result = spark.sql(""" + ALTER TABLE default.test_table OPTIMIZE INDEX idx_id + WITH (num_indices_to_merge = 0) + """).first() + + assert result.index_name == "idx_id" + assert result.fragments_indexed == before.num_unindexed_fragments + assert result.segments_after >= result.segments_before + + after = next( + row + for row in spark.sql("SHOW INDEXES IN default.test_table").collect() + if row.name == "idx_id" + ) + assert after.num_unindexed_fragments == 0 + def test_create_btree_index_on_nested_literal_dot_field(self, spark): """Test CREATE INDEX on nested struct fields, including literal dots.""" spark.sql(""" 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 e07fedd1e..fd6796f80 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 @@ -61,6 +61,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.concurrent.Callable; import java.util.concurrent.FutureTask; import java.util.stream.Collectors; @@ -202,6 +203,7 @@ public void commit(WriterCommitMessage[] messages) { .removedFragmentIds(removedFragmentIds) .updatedFragments(updatedFragments) .newFragments(newFragments) + .updateMode(Optional.of(Update.UpdateMode.RewriteRows)) .build(); CommitBuilder commitBuilder = diff --git a/lance-spark-3.4_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala b/lance-spark-3.4_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala index 06d605992..6b10418de 100644 --- a/lance-spark-3.4_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala +++ b/lance-spark-3.4_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala @@ -16,7 +16,7 @@ package org.apache.spark.sql.catalyst.parser.extensions import org.antlr.v4.runtime.ParserRuleContext import org.apache.spark.sql.catalyst.analysis.{UnresolvedIdentifier, UnresolvedRelation} import org.apache.spark.sql.catalyst.parser.{ParseException, ParserInterface} -import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} +import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceOptimizeIndex, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} import org.lance.spark.utils.{FieldPathUtils, ParserUtils} import java.util.Locale @@ -93,6 +93,19 @@ class LanceSqlExtensionsAstBuilder(delegate: ParserInterface) Optimize(table, args) } + override def visitOptimizeIndex(ctx: LanceSqlExtensionsParser.OptimizeIndexContext) + : LanceOptimizeIndex = { + val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) + val indexName = cleanIdentifier(ctx.indexName.getText) + val args = ctx.namedArgument().asScala.map(a => + LanceNamedArgument( + normalizedOptionName(a.identifier().getText), + a.constant().accept(this))) + .toSeq + + LanceOptimizeIndex(table, indexName, args) + } + override def visitVacuum(ctx: LanceSqlExtensionsParser.VacuumContext): Vacuum = { val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) val args = ctx.namedArgument().asScala.map(a => diff --git a/lance-spark-3.4_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java b/lance-spark-3.4_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java index 58aa0922f..e128ec056 100644 --- a/lance-spark-3.4_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java +++ b/lance-spark-3.4_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java @@ -22,6 +22,8 @@ import org.apache.spark.sql.catalyst.parser.extensions.LanceSqlExtensionsParser; import org.apache.spark.sql.catalyst.plans.logical.AddColumnsBackfill; import org.apache.spark.sql.catalyst.plans.logical.AddIndex; +import org.apache.spark.sql.catalyst.plans.logical.LanceNamedArgument; +import org.apache.spark.sql.catalyst.plans.logical.LanceOptimizeIndex; import org.apache.spark.sql.catalyst.plans.logical.Optimize; import org.apache.spark.sql.catalyst.plans.logical.ShowIndexes; import org.apache.spark.sql.catalyst.plans.logical.UpdateColumnsBackfill; @@ -182,6 +184,26 @@ public void testOptimizeNormalizesOptionNames() { assertEquals("target_rows_per_fragment", plan.args().apply(0).name()); } + @Test + public void testOptimizeIndexWithOptions() { + LanceSqlExtensionsParser parser = + createParser( + "ALTER TABLE `my-catalog`.`my-table` OPTIMIZE INDEX `my-idx` " + + "WITH (NUM_INDICES_TO_MERGE = 2)"); + LanceOptimizeIndex plan = + (LanceOptimizeIndex) astBuilder.visitSingleStatement(parser.singleStatement()); + + UnresolvedIdentifier table = (UnresolvedIdentifier) plan.table(); + assertEquals( + List.of("my-catalog", "my-table"), JavaConverters.seqAsJavaList(table.nameParts())); + assertEquals("my-idx", plan.indexName()); + + List args = JavaConverters.seqAsJavaList(plan.args()); + assertEquals(1, args.size()); + assertEquals("num_indices_to_merge", args.get(0).name()); + assertEquals(2L, args.get(0).value()); + } + @Test public void testShowIndexesWithBacktickedTableName() { LanceSqlExtensionsParser parser = createParser("SHOW INDEXES FROM `my-catalog`.`my-table`"); 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 c3921fd45..057f20b0a 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 @@ -66,6 +66,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.concurrent.Callable; import java.util.concurrent.FutureTask; import java.util.stream.Collectors; @@ -224,6 +225,7 @@ public void commit(WriterCommitMessage[] messages) { .removedFragmentIds(removedFragmentIds) .updatedFragments(updatedFragments) .newFragments(newFragments) + .updateMode(Optional.of(Update.UpdateMode.RewriteRows)) .build(); CommitBuilder commitBuilder = diff --git a/lance-spark-3.5_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala b/lance-spark-3.5_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala index 35a40208c..3ac6570ef 100644 --- a/lance-spark-3.5_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala +++ b/lance-spark-3.5_2.12/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala @@ -16,7 +16,7 @@ package org.apache.spark.sql.catalyst.parser.extensions import org.antlr.v4.runtime.ParserRuleContext import org.apache.spark.sql.catalyst.analysis.{UnresolvedIdentifier, UnresolvedRelation} import org.apache.spark.sql.catalyst.parser.{ParseException, ParserInterface} -import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} +import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceOptimizeIndex, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} import org.lance.spark.utils.{FieldPathUtils, ParserUtils} import java.util.Locale @@ -93,6 +93,19 @@ class LanceSqlExtensionsAstBuilder(delegate: ParserInterface) Optimize(table, args) } + override def visitOptimizeIndex(ctx: LanceSqlExtensionsParser.OptimizeIndexContext) + : LanceOptimizeIndex = { + val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) + val indexName = cleanIdentifier(ctx.indexName.getText) + val args = ctx.namedArgument().asScala.map(a => + LanceNamedArgument( + normalizedOptionName(a.identifier().getText), + a.constant().accept(this))) + .toSeq + + LanceOptimizeIndex(table, indexName, args) + } + override def visitVacuum(ctx: LanceSqlExtensionsParser.VacuumContext): Vacuum = { val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) val args = ctx.namedArgument().asScala.map(a => diff --git a/lance-spark-3.5_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java b/lance-spark-3.5_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java index 909ead579..8132483e8 100644 --- a/lance-spark-3.5_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java +++ b/lance-spark-3.5_2.12/src/test/java/org/lance/spark/update/LanceSqlExtensionsAstBuilderTest.java @@ -24,6 +24,8 @@ import org.apache.spark.sql.catalyst.plans.logical.AddIndex; import org.apache.spark.sql.catalyst.plans.logical.LanceCreateBranch; import org.apache.spark.sql.catalyst.plans.logical.LanceDropBranch; +import org.apache.spark.sql.catalyst.plans.logical.LanceNamedArgument; +import org.apache.spark.sql.catalyst.plans.logical.LanceOptimizeIndex; import org.apache.spark.sql.catalyst.plans.logical.LanceShowBranches; import org.apache.spark.sql.catalyst.plans.logical.Optimize; import org.apache.spark.sql.catalyst.plans.logical.ShowIndexes; @@ -206,6 +208,26 @@ public void testOptimizeNormalizesOptionNames() { assertEquals("target_rows_per_fragment", plan.args().apply(0).name()); } + @Test + public void testOptimizeIndexWithOptions() { + LanceSqlExtensionsParser parser = + createParser( + "ALTER TABLE `my-catalog`.`my-table` OPTIMIZE INDEX `my-idx` " + + "WITH (NUM_INDICES_TO_MERGE = 2)"); + LanceOptimizeIndex plan = + (LanceOptimizeIndex) astBuilder.visitSingleStatement(parser.singleStatement()); + + UnresolvedIdentifier table = (UnresolvedIdentifier) plan.table(); + assertEquals( + List.of("my-catalog", "my-table"), JavaConverters.seqAsJavaList(table.nameParts())); + assertEquals("my-idx", plan.indexName()); + + List args = JavaConverters.seqAsJavaList(plan.args()); + assertEquals(1, args.size()); + assertEquals("num_indices_to_merge", args.get(0).name()); + assertEquals(2L, args.get(0).value()); + } + @Test public void testShowIndexesWithBacktickedTableName() { LanceSqlExtensionsParser parser = createParser("SHOW INDEXES FROM `my-catalog`.`my-table`"); diff --git a/lance-spark-4.0_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala b/lance-spark-4.0_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala index c5b7ebf0d..b3a330d8b 100644 --- a/lance-spark-4.0_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala +++ b/lance-spark-4.0_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala @@ -16,7 +16,7 @@ package org.apache.spark.sql.catalyst.parser.extensions import org.antlr.v4.runtime.ParserRuleContext import org.apache.spark.sql.catalyst.analysis.{UnresolvedIdentifier, UnresolvedRelation} import org.apache.spark.sql.catalyst.parser.{ParseException, ParserInterface} -import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} +import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceOptimizeIndex, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} import org.lance.spark.utils.{FieldPathUtils, ParserUtils} import java.util.Locale @@ -93,6 +93,19 @@ class LanceSqlExtensionsAstBuilder(delegate: ParserInterface) Optimize(table, args) } + override def visitOptimizeIndex(ctx: LanceSqlExtensionsParser.OptimizeIndexContext) + : LanceOptimizeIndex = { + val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) + val indexName = cleanIdentifier(ctx.indexName.getText) + val args = ctx.namedArgument().asScala.map(a => + LanceNamedArgument( + normalizedOptionName(a.identifier().getText), + a.constant().accept(this))) + .toSeq + + LanceOptimizeIndex(table, indexName, args) + } + override def visitVacuum(ctx: LanceSqlExtensionsParser.VacuumContext): Vacuum = { val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) val args = ctx.namedArgument().asScala.map(a => diff --git a/lance-spark-4.1_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala b/lance-spark-4.1_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala index c5b7ebf0d..b3a330d8b 100644 --- a/lance-spark-4.1_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala +++ b/lance-spark-4.1_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala @@ -16,7 +16,7 @@ package org.apache.spark.sql.catalyst.parser.extensions import org.antlr.v4.runtime.ParserRuleContext import org.apache.spark.sql.catalyst.analysis.{UnresolvedIdentifier, UnresolvedRelation} import org.apache.spark.sql.catalyst.parser.{ParseException, ParserInterface} -import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} +import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceOptimizeIndex, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} import org.lance.spark.utils.{FieldPathUtils, ParserUtils} import java.util.Locale @@ -93,6 +93,19 @@ class LanceSqlExtensionsAstBuilder(delegate: ParserInterface) Optimize(table, args) } + override def visitOptimizeIndex(ctx: LanceSqlExtensionsParser.OptimizeIndexContext) + : LanceOptimizeIndex = { + val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) + val indexName = cleanIdentifier(ctx.indexName.getText) + val args = ctx.namedArgument().asScala.map(a => + LanceNamedArgument( + normalizedOptionName(a.identifier().getText), + a.constant().accept(this))) + .toSeq + + LanceOptimizeIndex(table, indexName, args) + } + override def visitVacuum(ctx: LanceSqlExtensionsParser.VacuumContext): Vacuum = { val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) val args = ctx.namedArgument().asScala.map(a => diff --git a/lance-spark-4.2_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala b/lance-spark-4.2_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala index c5b7ebf0d..b3a330d8b 100644 --- a/lance-spark-4.2_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala +++ b/lance-spark-4.2_2.13/src/main/scala/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensionsAstBuilder.scala @@ -16,7 +16,7 @@ package org.apache.spark.sql.catalyst.parser.extensions import org.antlr.v4.runtime.ParserRuleContext import org.apache.spark.sql.catalyst.analysis.{UnresolvedIdentifier, UnresolvedRelation} import org.apache.spark.sql.catalyst.parser.{ParseException, ParserInterface} -import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} +import org.apache.spark.sql.catalyst.plans.logical.{AddColumnsBackfill, AddIndex, LanceCreateBranch, LanceCreateTag, LanceDropBranch, LanceDropIndex, LanceDropTag, LanceNamedArgument, LanceOptimizeIndex, LanceShowBranches, LanceShowTags, LogicalPlan, Optimize, SetUnenforcedPrimaryKey, ShowIndexes, UpdateColumnsBackfill, Vacuum} import org.lance.spark.utils.{FieldPathUtils, ParserUtils} import java.util.Locale @@ -93,6 +93,19 @@ class LanceSqlExtensionsAstBuilder(delegate: ParserInterface) Optimize(table, args) } + override def visitOptimizeIndex(ctx: LanceSqlExtensionsParser.OptimizeIndexContext) + : LanceOptimizeIndex = { + val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) + val indexName = cleanIdentifier(ctx.indexName.getText) + val args = ctx.namedArgument().asScala.map(a => + LanceNamedArgument( + normalizedOptionName(a.identifier().getText), + a.constant().accept(this))) + .toSeq + + LanceOptimizeIndex(table, indexName, args) + } + override def visitVacuum(ctx: LanceSqlExtensionsParser.VacuumContext): Vacuum = { val table = UnresolvedIdentifier(visitMultipartIdentifier(ctx.multipartIdentifier())) val args = ctx.namedArgument().asScala.map(a => diff --git a/lance-spark-base_2.12/src/main/antlr4/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensions.g4 b/lance-spark-base_2.12/src/main/antlr4/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensions.g4 index 788525249..adc4df35a 100644 --- a/lance-spark-base_2.12/src/main/antlr4/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensions.g4 +++ b/lance-spark-base_2.12/src/main/antlr4/org/apache/spark/sql/catalyst/parser/extensions/LanceSqlExtensions.g4 @@ -23,6 +23,8 @@ statement | ALTER TABLE multipartIdentifier UPDATE COLUMNS columnList FROM identifier #updateColumnsBackfill | ALTER TABLE multipartIdentifier CREATE INDEX indexName=identifier USING method=identifier '(' fieldPathList ')' (WITH '(' (namedArgument (',' namedArgument)*)? ')')? #createIndex | ALTER TABLE multipartIdentifier DROP INDEX indexName=identifier #dropIndex + | ALTER TABLE multipartIdentifier OPTIMIZE INDEX indexName=identifier + (WITH '(' (namedArgument (',' namedArgument)*)? ')')? #optimizeIndex | ALTER TABLE multipartIdentifier CREATE BRANCH (IF NOT EXISTS)? branchName=identifier (AS OF VERSION refMainVersion=versionNumber)? #createBranchRefMain | ALTER TABLE multipartIdentifier CREATE BRANCH (IF NOT EXISTS)? branchName=identifier @@ -172,4 +174,3 @@ fragment DIGIT fragment LETTER : [A-Z] ; - diff --git a/lance-spark-base_2.12/src/main/java/org/lance/spark/search/LanceSearchColumnarPartitionReader.java b/lance-spark-base_2.12/src/main/java/org/lance/spark/search/LanceSearchColumnarPartitionReader.java index 80c30a344..3ed59586d 100644 --- a/lance-spark-base_2.12/src/main/java/org/lance/spark/search/LanceSearchColumnarPartitionReader.java +++ b/lance-spark-base_2.12/src/main/java/org/lance/spark/search/LanceSearchColumnarPartitionReader.java @@ -84,7 +84,7 @@ private void openArrowReader() throws IOException { throw new IOException("Lance namespace is required for search"); } try { - byte[] bytes = namespace.queryTable(query.toQueryTableRequest()); + byte[] bytes = namespace.queryTable(query.toQueryTableRequest()).getData(); arrowReader = new ArrowFileReader( new ByteArrayReadableSeekableByteChannel(bytes), LanceRuntime.allocator()); diff --git a/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/OptimizeIndex.scala b/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/OptimizeIndex.scala new file mode 100644 index 000000000..80ff8a904 --- /dev/null +++ b/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/OptimizeIndex.scala @@ -0,0 +1,45 @@ +/* + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.spark.sql.catalyst.plans.logical + +import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeReference} +import org.apache.spark.sql.types.{DataTypes, StructField, StructType} + +/** Logical plan for maintaining an existing named Lance index. */ +case class LanceOptimizeIndex( + table: LogicalPlan, + indexName: String, + args: Seq[LanceNamedArgument]) extends Command { + + override def children: Seq[LogicalPlan] = Seq(table) + + override def output: Seq[Attribute] = LanceOptimizeIndexOutputType.SCHEMA + + override def simpleString(maxFields: Int): String = s"LanceOptimizeIndex($indexName)" + + override protected def withNewChildrenInternal(newChildren: IndexedSeq[LogicalPlan]) + : LanceOptimizeIndex = { + copy(table = newChildren(0)) + } +} + +object LanceOptimizeIndexOutputType { + val SCHEMA: Seq[Attribute] = StructType( + Array( + StructField("index_name", DataTypes.StringType, nullable = false), + StructField("fragments_indexed", DataTypes.LongType, nullable = false), + StructField("segments_before", DataTypes.LongType, nullable = false), + StructField("segments_after", DataTypes.LongType, nullable = false))) + .map(field => AttributeReference(field.name, field.dataType, field.nullable, field.metadata)()) +} diff --git a/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/execution/datasources/v2/LanceDataSourceV2Strategy.scala b/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/execution/datasources/v2/LanceDataSourceV2Strategy.scala index f17e0d52c..142d9a651 100644 --- a/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/execution/datasources/v2/LanceDataSourceV2Strategy.scala +++ b/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/execution/datasources/v2/LanceDataSourceV2Strategy.scala @@ -20,6 +20,8 @@ import org.apache.spark.sql.catalyst.plans.logical._ import org.apache.spark.sql.connector.catalog._ import org.apache.spark.sql.execution.{SparkPlan, SparkStrategy} +import java.util.Locale + case class LanceDataSourceV2Strategy(session: SparkSession) extends SparkStrategy with PredicateHelper { @@ -51,6 +53,13 @@ case class LanceDataSourceV2Strategy(session: SparkSession) extends SparkStrateg case LanceDropIndex(ResolvedIdentifier(catalog, ident), indexName) => LanceDropIndexExec(asTableCatalog(catalog), ident, indexName.toLowerCase) :: Nil + case LanceOptimizeIndex(ResolvedIdentifier(catalog, ident), indexName, args) => + LanceOptimizeIndexExec( + asTableCatalog(catalog), + ident, + indexName.toLowerCase(Locale.ROOT), + args) :: Nil + case LanceCreateBranch(ResolvedIdentifier(catalog, ident), branchName, ref, ifNotExists) => LanceCreateBranchExec(asTableCatalog(catalog), ident, branchName, ref, ifNotExists) :: Nil diff --git a/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/execution/datasources/v2/OptimizeIndexExec.scala b/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/execution/datasources/v2/OptimizeIndexExec.scala new file mode 100644 index 000000000..189de8931 --- /dev/null +++ b/lance-spark-base_2.12/src/main/scala/org/apache/spark/sql/execution/datasources/v2/OptimizeIndexExec.scala @@ -0,0 +1,188 @@ +/* + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.spark.sql.execution.datasources.v2 + +import org.apache.spark.sql.catalyst.InternalRow +import org.apache.spark.sql.catalyst.expressions.{Attribute, GenericInternalRow} +import org.apache.spark.sql.catalyst.plans.logical.{LanceNamedArgument, LanceOptimizeIndexOutputType} +import org.apache.spark.sql.connector.catalog.{Identifier, TableCatalog} +import org.apache.spark.unsafe.types.UTF8String +import org.lance.Dataset +import org.lance.index.{Index, IndexCriteria, OptimizeOptions} +import org.lance.operation.CreateIndex +import org.lance.spark.LanceDataset +import org.lance.spark.utils.Utils + +import java.util.{Collections, Locale} + +import scala.collection.JavaConverters._ + +object LanceOptimizeIndexExec { + private val SystemIndexNames = Set("__lance_frag_reuse", "__lance_mem_wal") + + private def isSystemIndex(indexName: String): Boolean = + indexName != null && SystemIndexNames.exists(_.equalsIgnoreCase(indexName)) +} + +/** Driver-side execution of ALTER TABLE ... OPTIMIZE INDEX. */ +case class LanceOptimizeIndexExec( + catalog: TableCatalog, + ident: Identifier, + indexName: String, + args: Seq[LanceNamedArgument]) extends LeafV2CommandExec { + + override def output: Seq[Attribute] = LanceOptimizeIndexOutputType.SCHEMA + + private case class IndexState(segmentCount: Long) + + private case class IndexDelta( + fragmentsIndexed: Long, + segmentsAdded: Long, + segmentsRemoved: Long) + + private def buildOptions(): OptimizeOptions = { + val normalizedArgs = args.map(arg => arg.name.toLowerCase(Locale.ROOT) -> arg) + val duplicateArgs = normalizedArgs.groupBy(_._1).collect { + case (name, values) if values.size > 1 => name + }.toSeq.sorted + if (duplicateArgs.nonEmpty) { + throw new IllegalArgumentException( + s"Duplicate OPTIMIZE INDEX options: ${duplicateArgs.mkString(", ")}") + } + + val argsMap = normalizedArgs.toMap + val supported = Set("num_indices_to_merge") + val unsupported = argsMap.keySet.diff(supported).toSeq.sorted + if (unsupported.nonEmpty) { + throw new IllegalArgumentException( + s"Unsupported OPTIMIZE INDEX options: ${unsupported.mkString(", ")}") + } + + val builder = OptimizeOptions.builder() + .indexNames(Collections.singletonList(indexName)) + + argsMap.get("num_indices_to_merge").foreach { arg => + val value = arg.value match { + case number: java.lang.Long => number.longValue() + case other => + throw new IllegalArgumentException( + s"num_indices_to_merge must be a non-negative integer, got: $other") + } + if (value < 0 || value > Int.MaxValue) { + throw new IllegalArgumentException( + s"num_indices_to_merge must be between 0 and ${Int.MaxValue}, got: $value") + } + builder.numIndicesToMerge(value.toInt) + } + + builder.build() + } + + private def indexState(dataset: Dataset): IndexState = { + val criteria = new IndexCriteria.Builder().hasName(indexName).build() + val description = dataset.describeIndices(criteria).asScala + .find(_.getName == indexName) + .getOrElse(throw new IllegalArgumentException(s"Index '$indexName' does not exist")) + IndexState(description.getSegments.size().toLong) + } + + private def requiredFragments(index: Index): Set[Integer] = { + val fragments = index.fragments() + if (!fragments.isPresent) { + throw new IllegalStateException( + s"Lance index segment '${index.uuid()}' for '$indexName' has no fragment metadata") + } + fragments.get().asScala.toSet + } + + private def indexDelta(dataset: Dataset, beforeVersion: Long, afterVersion: Long): IndexDelta = { + if (afterVersion <= beforeVersion) { + throw new IllegalStateException( + s"OPTIMIZE INDEX '$indexName' produced invalid dataset version change: " + + s"$beforeVersion -> $afterVersion") + } + + val maybeTransaction = dataset.readTransaction() + if (!maybeTransaction.isPresent) { + throw new IllegalStateException( + s"Dataset version $afterVersion has no transaction for OPTIMIZE INDEX '$indexName'") + } + + val transaction = maybeTransaction.get() + try { + transaction.operation() match { + case operation: CreateIndex => + val added = operation.getNewIndices.asScala.filter(_.name() == indexName) + val removed = operation.getRemovedIndices.asScala.filter(_.name() == indexName) + val addedFragments = added.flatMap(requiredFragments).toSet + val removedFragments = removed.flatMap(requiredFragments).toSet + IndexDelta( + addedFragments.diff(removedFragments).size.toLong, + added.size.toLong, + removed.size.toLong) + case operation => + throw new IllegalStateException( + s"Dataset version $afterVersion was created by '${operation.name()}', " + + s"not OPTIMIZE INDEX '$indexName'") + } + } finally { + transaction.close() + } + } + + override protected def run(): Seq[InternalRow] = { + val lanceDataset = LanceDataset.requireWritable(catalog.loadTable(ident), "OptimizeIndex") + if (LanceOptimizeIndexExec.isSystemIndex(indexName)) { + throw new IllegalArgumentException(s"Cannot optimize system index '$indexName'") + } + val options = buildOptions() + + val dataset = Utils.openDatasetBuilder(lanceDataset.readOptions()) + .initialStorageOptions(lanceDataset.getInitialStorageOptions) + .build() + try { + val before = indexState(dataset) + val beforeVersion = dataset.version() + dataset.optimizeIndices(options) + val afterVersion = dataset.version() + val after = indexState(dataset) + + val delta = if (afterVersion == beforeVersion) { + if (after != before) { + throw new IllegalStateException( + s"Index '$indexName' changed without a new dataset version") + } + IndexDelta(0L, 0L, 0L) + } else { + indexDelta(dataset, beforeVersion, afterVersion) + } + val expectedSegmentsAfter = + before.segmentCount - delta.segmentsRemoved + delta.segmentsAdded + if (expectedSegmentsAfter != after.segmentCount) { + throw new IllegalStateException( + s"OPTIMIZE INDEX '$indexName' reported an inconsistent segment change: " + + s"${before.segmentCount} - ${delta.segmentsRemoved} + ${delta.segmentsAdded} != " + + s"${after.segmentCount}") + } + + Seq(new GenericInternalRow(Array[Any]( + UTF8String.fromString(indexName), + delta.fragmentsIndexed, + before.segmentCount, + after.segmentCount))) + } finally { + dataset.close() + } + } +} diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/LanceRuntimeQueryTableSupportTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/LanceRuntimeQueryTableSupportTest.java index cc961e552..ca68de025 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/LanceRuntimeQueryTableSupportTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/LanceRuntimeQueryTableSupportTest.java @@ -15,6 +15,7 @@ import org.lance.namespace.LanceNamespace; import org.lance.namespace.model.QueryTableRequest; +import org.lance.namespace.model.QueryTableResponse; import org.apache.arrow.memory.BufferAllocator; import org.junit.jupiter.api.Test; @@ -97,8 +98,8 @@ public static class QueryingNamespace extends CatalogOnlyNamespace { public QueryingNamespace() {} @Override - public byte[] queryTable(QueryTableRequest request) { - return new byte[0]; + public QueryTableResponse queryTable(QueryTableRequest request) { + return new QueryTableResponse().data(new byte[0]); } } } diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/branch/BaseBranchDDLTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/branch/BaseBranchDDLTest.java index 22d202a21..5e863674c 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/branch/BaseBranchDDLTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/branch/BaseBranchDDLTest.java @@ -665,7 +665,7 @@ public void testInsertIntoBranchIdentifierFails() { } @Test - public void testBranchIdentifierRejectsCreateIndexAndVacuum() throws Exception { + public void testBranchIdentifierRejectsIndexMaintenanceAndVacuum() throws Exception { prepareDatasetWithHistory(); spark.sql(String.format("alter table %s create branch audit", fullTable)); @@ -694,6 +694,16 @@ public void testBranchIdentifierRejectsCreateIndexAndVacuum() throws Exception { .collectAsList()); Assertions.assertTrue(exceptionChainMessages(createIndex).contains("Writes are not supported")); + Exception optimizeIndex = + Assertions.assertThrows( + Exception.class, + () -> + spark + .sql("alter table " + fullTable + ".branch_audit optimize index id_idx") + .collectAsList()); + Assertions.assertTrue( + exceptionChainMessages(optimizeIndex).contains("Writes are not supported")); + Exception vacuum = Assertions.assertThrows( Exception.class, diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/read/BaseFtsCatalogOnlyNamespaceTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/read/BaseFtsCatalogOnlyNamespaceTest.java index 9cb54b0c7..f4ca55539 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/read/BaseFtsCatalogOnlyNamespaceTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/read/BaseFtsCatalogOnlyNamespaceTest.java @@ -34,11 +34,13 @@ import org.lance.namespace.model.ListTablesRequest; import org.lance.namespace.model.ListTablesResponse; import org.lance.namespace.model.NamespaceExistsRequest; +import org.lance.namespace.model.NamespaceExistsResponse; import org.lance.namespace.model.RegisterTableRequest; import org.lance.namespace.model.RegisterTableResponse; import org.lance.namespace.model.RenameTableRequest; import org.lance.namespace.model.RenameTableResponse; import org.lance.namespace.model.TableExistsRequest; +import org.lance.namespace.model.TableExistsResponse; import org.apache.arrow.memory.BufferAllocator; import org.apache.spark.sql.Dataset; @@ -276,8 +278,8 @@ public DropNamespaceResponse dropNamespace(DropNamespaceRequest request) { } @Override - public void namespaceExists(NamespaceExistsRequest request) { - delegate.namespaceExists(request); + public NamespaceExistsResponse namespaceExists(NamespaceExistsRequest request) { + return delegate.namespaceExists(request); } @Override @@ -306,8 +308,8 @@ public DeregisterTableResponse deregisterTable(DeregisterTableRequest request) { } @Override - public void tableExists(TableExistsRequest request) { - delegate.tableExists(request); + public TableExistsResponse tableExists(TableExistsRequest request) { + return delegate.tableExists(request); } @Override diff --git a/lance-spark-base_2.12/src/test/java/org/lance/spark/update/BaseOptimizeTest.java b/lance-spark-base_2.12/src/test/java/org/lance/spark/update/BaseOptimizeTest.java index f618d328b..b63d29dd7 100644 --- a/lance-spark-base_2.12/src/test/java/org/lance/spark/update/BaseOptimizeTest.java +++ b/lance-spark-base_2.12/src/test/java/org/lance/spark/update/BaseOptimizeTest.java @@ -13,17 +13,29 @@ */ package org.lance.spark.update; +import org.lance.Transaction; +import org.lance.operation.CreateIndex; +import org.lance.spark.LanceDataset; +import org.lance.spark.utils.Utils; + import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.catalyst.plans.logical.LanceNamedArgument; +import org.apache.spark.sql.connector.catalog.Identifier; +import org.apache.spark.sql.connector.catalog.TableCatalog; +import org.apache.spark.sql.execution.datasources.v2.LanceOptimizeIndexExec; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; +import scala.collection.JavaConverters; import java.io.IOException; +import java.lang.reflect.Method; import java.nio.file.Path; +import java.util.Collections; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -137,4 +149,177 @@ public void testWithoutArgs() { Assertions.assertEquals("[10,1,10,1]", result.collectAsList().get(0).toString()); } + + @Test + public void testOptimizeIndex() { + prepareDataset(); + spark.sql(String.format("alter table %s create index idx_id using zonemap (id)", fullTable)); + spark.sql(String.format("insert into %s values (10, 'text_10')", fullTable)); + + Row before = + spark.sql(String.format("show indexes in %s", fullTable)).collectAsList().stream() + .filter(row -> "idx_id".equals(row.getAs("name"))) + .findFirst() + .orElseThrow(() -> new AssertionError("Index not found: idx_id")); + long unindexedFragments = before.getAs("num_unindexed_fragments"); + Assertions.assertTrue(unindexedFragments > 0); + + Dataset result = + spark.sql(String.format("alter table %s optimize index idx_id", fullTable)); + + Assertions.assertEquals( + "StructType(StructField(index_name,StringType,false),StructField(fragments_indexed,LongType,false),StructField(segments_before,LongType,false),StructField(segments_after,LongType,false))", + result.schema().toString()); + Row optimized = result.collectAsList().get(0); + Assertions.assertEquals("idx_id", optimized.getAs("index_name")); + Assertions.assertEquals(unindexedFragments, optimized.getAs("fragments_indexed")); + Assertions.assertTrue(optimized.getAs("segments_before") > 0); + Assertions.assertTrue(optimized.getAs("segments_after") > 0); + + Row after = + spark.sql(String.format("show indexes in %s", fullTable)).collectAsList().stream() + .filter(row -> "idx_id".equals(row.getAs("name"))) + .findFirst() + .orElseThrow(() -> new AssertionError("Index not found: idx_id")); + Assertions.assertEquals(0L, after.getAs("num_unindexed_fragments")); + + Row noOp = + spark + .sql(String.format("alter table %s optimize index idx_id", fullTable)) + .collectAsList() + .get(0); + Assertions.assertEquals(0L, noOp.getAs("fragments_indexed")); + Assertions.assertEquals( + noOp.getAs("segments_before").longValue(), + noOp.getAs("segments_after").longValue()); + } + + @Test + public void testOptimizeIndexAcceptsMemWalCatchUpCommit() throws Exception { + String protocolTableName = tableName + "_protocol"; + String protocolTable = catalogName + ".default." + protocolTableName; + spark.sql( + String.format( + "create table %s (id int, text string) using lance " + + "partitioned by (bucket(4, text))", + protocolTable)); + spark.sql(String.format("insert into %s values (1, 'same-shard')", protocolTable)); + spark.sql( + String.format("alter table %s create index idx_id using zonemap (id)", protocolTable)); + + Row noOp = + spark + .sql(String.format("alter table %s optimize index idx_id", protocolTable)) + .collectAsList() + .get(0); + Assertions.assertEquals(0L, noOp.getAs("fragments_indexed")); + + try (org.lance.Dataset dataset = openDataset(protocolTableName)) { + long beforeVersion = dataset.version(); + CreateIndex emptyCreateIndex = + CreateIndex.builder() + .withNewIndices(Collections.emptyList()) + .withRemovedIndices(Collections.emptyList()) + .build(); + try (Transaction transaction = new Transaction(beforeVersion, emptyCreateIndex); + org.lance.Dataset committed = dataset.commitTransaction(transaction)) { + Assertions.assertEquals(beforeVersion + 1, committed.version()); + + LanceOptimizeIndexExec exec = + new LanceOptimizeIndexExec( + null, + null, + "idx_id", + JavaConverters.asScalaBuffer(Collections.emptyList()).toSeq()); + Method indexDelta = + LanceOptimizeIndexExec.class.getDeclaredMethod( + "indexDelta", org.lance.Dataset.class, long.class, long.class); + indexDelta.setAccessible(true); + Object delta = indexDelta.invoke(exec, committed, beforeVersion, committed.version()); + + Assertions.assertEquals(0L, metric(delta, "fragmentsIndexed")); + Assertions.assertEquals(0L, metric(delta, "segmentsAdded")); + Assertions.assertEquals(0L, metric(delta, "segmentsRemoved")); + } + } + } + + private org.lance.Dataset openDataset(String currentTableName) throws Exception { + TableCatalog catalog = + (TableCatalog) spark.sessionState().catalogManager().catalog(catalogName); + LanceDataset table = + (LanceDataset) catalog.loadTable(Identifier.of(new String[] {"default"}, currentTableName)); + return Utils.openDatasetBuilder(table.readOptions()) + .initialStorageOptions(table.getInitialStorageOptions()) + .build(); + } + + private static long metric(Object delta, String name) throws Exception { + return ((Long) delta.getClass().getDeclaredMethod(name).invoke(delta)).longValue(); + } + + @Test + public void testOptimizeMissingIndex() { + prepareDataset(); + + assertOptimizeIndexFails("missing", "", "Index 'missing' does not exist"); + } + + @Test + public void testOptimizeRejectsSystemIndex() { + prepareDataset(); + + assertOptimizeIndexFails( + "__lance_frag_reuse", "", "Cannot optimize system index '__lance_frag_reuse'"); + } + + @Test + public void testOptimizeIndexValidatesOptions() { + prepareDataset(); + spark.sql(String.format("alter table %s create index idx_id using zonemap (id)", fullTable)); + + assertOptimizeIndexFails( + "idx_id", "with (unknown_option=1)", "Unsupported OPTIMIZE INDEX options: unknown_option"); + assertOptimizeIndexFails( + "idx_id", + "with (num_indices_to_merge=0, NUM_INDICES_TO_MERGE=1)", + "Duplicate OPTIMIZE INDEX options: num_indices_to_merge"); + assertOptimizeIndexFails( + "idx_id", + "with (num_indices_to_merge='two')", + "num_indices_to_merge must be a non-negative integer"); + assertOptimizeIndexFails( + "idx_id", + "with (num_indices_to_merge=-1)", + "num_indices_to_merge must be between 0 and 2147483647"); + assertOptimizeIndexFails( + "idx_id", + "with (num_indices_to_merge=2147483648)", + "num_indices_to_merge must be between 0 and 2147483647"); + } + + private void assertOptimizeIndexFails(String indexName, String options, String expectedMessage) { + Exception error = + Assertions.assertThrows( + Exception.class, + () -> + spark + .sql( + String.format( + "alter table %s optimize index %s %s", fullTable, indexName, options)) + .collectAsList()); + Assertions.assertTrue( + exceptionChainMessages(error).contains(expectedMessage), + () -> "Expected error containing '" + expectedMessage + "', got: " + error); + } + + private static String exceptionChainMessages(Throwable throwable) { + StringBuilder messages = new StringBuilder(); + for (Throwable current = throwable; current != null; current = current.getCause()) { + if (current.getMessage() != null) { + messages.append(current.getMessage()).append('\n'); + } + } + return messages.toString(); + } } diff --git a/pom.xml b/pom.xml index 91d9e4fee..67c4bedd1 100644 --- a/pom.xml +++ b/pom.xml @@ -51,8 +51,8 @@ 0.8.0-beta.1 - 11.0.0-beta.21 - 0.8.6 + 12.0.0-beta.14 + 0.11.1 0.4.0 14.0.2