This is an automated email from the ASF dual-hosted git repository.
FANNG1 pushed a commit to branch fix/lance-arrow-field-metadata
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/fix/lance-arrow-field-metadata
by this push:
new 165f71a824 [#13338] fix(lance): Expose lance.blob and pick the Lance
file format version for blobs
165f71a824 is described below
commit 165f71a8245ef8e9f8fdfaf282f2cd4879212230
Author: fanng <[email protected]>
AuthorDate: Fri Sep 25 12:51:36 2026 +0900
[#13338] fix(lance): Expose lance.blob and pick the Lance file format
version for blobs
---
.../lakehouse/lance/LanceTableOperations.java | 27 +-
.../lakehouse/lance/TestLanceTableOperations.java | 74 +++++
docs/lakehouse-generic-lance-table.md | 44 +--
.../lance/common/ops/gravitino/LanceBlobTypes.java | 344 ++++++++------------
.../ops/gravitino/LanceDataTypeConverter.java | 26 ++
.../common/ops/gravitino/TestLanceBlobTypes.java | 349 ++++++++++-----------
.../ops/gravitino/TestLanceDataTypeConverter.java | 84 +++--
7 files changed, 487 insertions(+), 461 deletions(-)
diff --git
a/catalogs/catalog-lakehouse-generic/src/main/java/org/apache/gravitino/catalog/lakehouse/lance/LanceTableOperations.java
b/catalogs/catalog-lakehouse-generic/src/main/java/org/apache/gravitino/catalog/lakehouse/lance/LanceTableOperations.java
index 4d9c3dceda..1ed701113b 100644
---
a/catalogs/catalog-lakehouse-generic/src/main/java/org/apache/gravitino/catalog/lakehouse/lance/LanceTableOperations.java
+++
b/catalogs/catalog-lakehouse-generic/src/main/java/org/apache/gravitino/catalog/lakehouse/lance/LanceTableOperations.java
@@ -72,6 +72,7 @@ import org.apache.gravitino.utils.ExceptionMessages;
import org.apache.gravitino.utils.PrincipalUtils;
import org.lance.Dataset;
import org.lance.ReadOptions;
+import org.lance.WriteDatasetBuilder;
import org.lance.WriteParams;
import org.lance.index.DistanceType;
import org.lance.index.IndexOptions;
@@ -419,13 +420,7 @@ public class LanceTableOperations extends
ManagedTableOperations {
Map<String, String> storageProps =
LancePropertiesUtils.resolveLanceStorageOptions(catalogProperties,
properties);
- try (Dataset ignored =
- Dataset.write()
- .schema(convertColumnsToArrowSchema(columns))
- .uri(location)
- .mode(WriteParams.WriteMode.CREATE)
- .storageOptions(storageProps)
- .execute()) {
+ try (Dataset ignored = createDataset(location, columns, storageProps)) {
// Only create the table metadata in Gravitino after the Lance dataset
is successfully
// created.
long datasetVersion = ignored.version();
@@ -461,6 +456,23 @@ public class LanceTableOperations extends
ManagedTableOperations {
}
}
+ /**
+ * Creates an empty Lance dataset for the columns. The Lance file format
version is chosen from
+ * the blob columns, because blob v2 and legacy blob each require a
different version.
+ */
+ Dataset createDataset(String location, Column[] columns, Map<String, String>
storageOptions) {
+ Schema schema = convertColumnsToArrowSchema(columns);
+ WriteDatasetBuilder builder =
+ Dataset.write()
+ .schema(schema)
+ .uri(location)
+ .mode(WriteParams.WriteMode.CREATE)
+ .storageOptions(storageOptions);
+ LanceDataTypeConverter.requiredFileFormatVersion(schema.getFields())
+ .ifPresent(builder::dataStorageVersion);
+ return builder.execute();
+ }
+
private Schema convertColumnsToArrowSchema(Column[] columns) {
List<Field> fields =
Arrays.stream(columns)
@@ -854,6 +866,7 @@ public class LanceTableOperations extends
ManagedTableOperations {
Field field =
LanceDataTypeConverter.CONVERTER.toArrowField(
columnName, addColumn.getDataType(), true);
+ LanceDataTypeConverter.checkFileFormatVersion(field,
dataset.getLanceFileFormatVersion());
dataset.addColumns(List.of(field));
} else if (change instanceof TableChange.DeleteColumn deleteColumn) {
dataset.dropColumns(List.of(String.join(".",
deleteColumn.fieldName())));
diff --git
a/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/TestLanceTableOperations.java
b/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/TestLanceTableOperations.java
index 890c718a4d..6d688b0d7b 100644
---
a/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/TestLanceTableOperations.java
+++
b/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/TestLanceTableOperations.java
@@ -50,6 +50,7 @@ import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.UserPrincipal;
import org.apache.gravitino.catalog.ManagedSchemaOperations;
import org.apache.gravitino.exceptions.OptimisticLockException;
+import org.apache.gravitino.lance.common.ops.gravitino.LanceDataTypeConverter;
import org.apache.gravitino.meta.AuditInfo;
import org.apache.gravitino.meta.ColumnEntity;
import org.apache.gravitino.meta.TableEntity;
@@ -68,6 +69,7 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.lance.Dataset;
import org.lance.Version;
@@ -898,6 +900,78 @@ public class TestLanceTableOperations {
Mockito.verify(dataset).getVersion();
}
+ @ParameterizedTest
+ @CsvSource({"lance.blob, 2.2", "lance.blob.legacy, 2.1"})
+ public void testCreateDatasetWithBlobColumnRoundTrips(String blobType,
String fileFormatVersion) {
+ String location = tempDir.resolve("blob-" + fileFormatVersion).toString();
+ Column[] columns =
+ new Column[] {
+ Column.of("id", Types.IntegerType.get()),
+ Column.of("image", Types.ExternalType.of(blobType))
+ };
+
+ try (Dataset dataset = lanceTableOps.createDataset(location, columns,
Map.of())) {
+ Assertions.assertEquals(fileFormatVersion,
dataset.getLanceFileFormatVersion());
+ }
+ try (Dataset dataset = lanceTableOps.openDataset(location, Map.of())) {
+ Assertions.assertEquals(
+ Types.ExternalType.of(blobType),
+
LanceDataTypeConverter.CONVERTER.toGravitino(dataset.getSchema().findField("image")));
+ }
+ }
+
+ @Test
+ public void testCreateDatasetRejectsMixedBlobColumns() {
+ Column[] columns =
+ new Column[] {
+ Column.of("image", Types.ExternalType.of("lance.blob")),
+ Column.of("legacy_image", Types.ExternalType.of("lance.blob.legacy"))
+ };
+
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ lanceTableOps.createDataset(
+ tempDir.resolve("mixed-blob").toString(), columns, Map.of()));
+ }
+
+ @Test
+ public void testAddBlobColumnChecksFileFormatVersion() {
+ String location = tempDir.resolve("add-blob").toString();
+ try (Dataset ignored =
+ lanceTableOps.createDataset(
+ location, new Column[] {Column.of("id", Types.IntegerType.get())},
Map.of())) {
+ // A dataset without blob columns uses Lance's default file format
version.
+ }
+ Table table = mock(Table.class);
+ when(table.properties()).thenReturn(Map.of(Table.PROPERTY_LOCATION,
location));
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ lanceTableOps.handleLanceTableChange(
+ table,
+ new TableChange[] {
+ TableChange.addColumn(
+ new String[] {"image"},
Types.ExternalType.of("lance.blob"))
+ }));
+ Assertions.assertTrue(exception.getMessage().contains("2.2"),
exception.getMessage());
+
+ lanceTableOps.handleLanceTableChange(
+ table,
+ new TableChange[] {
+ TableChange.addColumn(
+ new String[] {"legacy_image"},
Types.ExternalType.of("lance.blob.legacy"))
+ });
+ try (Dataset dataset = lanceTableOps.openDataset(location, Map.of())) {
+ Assertions.assertEquals(
+ Types.ExternalType.of("lance.blob.legacy"),
+ LanceDataTypeConverter.CONVERTER.toGravitino(
+ dataset.getSchema().findField("legacy_image")));
+ }
+ }
+
@Test
public void
testVersionCheckSkipsSchemaReadForEmptySchemaWhenVersionUnchanged() throws
Exception {
lanceTableOps.setCatalogProperties(Map.of(LANCE_SCHEMA_REFRESH_MODE,
"version-check"));
diff --git a/docs/lakehouse-generic-lance-table.md
b/docs/lakehouse-generic-lance-table.md
index 4252c23c58..9740a2cf89 100644
--- a/docs/lakehouse-generic-lance-table.md
+++ b/docs/lakehouse-generic-lance-table.md
@@ -75,8 +75,8 @@ Lance uses Apache Arrow for table schemas. The following
table shows type mappin
| `Interval_year` | Not supported by Lance |
| `Interval_day` | `Duration(Microsecond)` |
| `External(arrow_field_json_str)` | Any Arrow Field |
-| `External("lance.blob.v1")` | Lance legacy blob |
-| `External("lance.blob.v2(...)")` | Lance blob v2 |
+| `External("lance.blob")` | Lance blob |
+| `External("lance.blob.legacy")` | Lance legacy blob |
### External Types
@@ -102,34 +102,38 @@ For Arrow types not natively mapped in Gravitino, use the
`External(arrow_field_
Gravitino types cannot carry Arrow field metadata. When loading a Lance table,
only Lance blob metadata
is recognized (see [Blob Types](#blob-types)); other field metadata is
ignored. Blob fields nested in a
`Struct`, `List`, `Map` or `Union` keep the container as a native Gravitino
type, for example
-`List(External("lance.blob.v2"))`. As with other list columns, the list child
is written back as
+`List(External("lance.blob"))`. As with other list columns, the list child is
written back as
`element`, whatever name it had in the Lance dataset (Lance and pyarrow use
`item`).
### Blob Types
Lance blob columns use a readable external type instead of Arrow JSON:
-| External Type Definition | Arrow Field
|
-|----------------------------------------------|---------------------------------------------------------------------------------------------------------|
-| `External("lance.blob.v1")` | `LargeBinary` with metadata
`lance-encoding:blob=true` |
-| `External("lance.blob.v2")` | `Struct<data: LargeBinary,
uri: Utf8>` with metadata `ARROW:extension:name=lance.blob.v2` |
-| `External("lance.blob.v2(with_range=true)")` | `Struct<data: LargeBinary,
uri: Utf8, position: UInt64, size: UInt64>` with the same extension metadata |
+| External Type Definition | Arrow Field
| Lance File Format Version |
+|-----------------------------------|-------------------------------------------------------------------------------------------|---------------------------|
+| `External("lance.blob")` | `Struct<data: LargeBinary, uri: Utf8>`
with metadata `ARROW:extension:name=lance.blob.v2` | 2.2 or later |
+| `External("lance.blob.legacy")` | `LargeBinary` with metadata
`lance-encoding:blob=true` | 2.1 or earlier
|
-`lance.blob.v2` accepts these optional parameters, written as
`lance.blob.v2(key=value, ...)`:
+Use `lance.blob` for new tables. `lance.blob.legacy` is only for datasets that
already use Lance's
+legacy blob encoding.
-| Parameter | Arrow Field Metadata
| Value |
-|----------------------------|------------------------------------------------|-------------------------------------------------------------|
-| `with_range` | -
| `true` or `false` (default), adds `position` and `size` |
-| `inline_size_threshold` | `lance-encoding:blob-inline-size-threshold`
| Integer >= 0 |
-| `dedicated_size_threshold` | `lance-encoding:blob-dedicated-size-threshold`
| Integer > 0 |
-| `pack_file_size_threshold` | `lance-encoding:blob-pack-file-size-threshold`
| Integer > 0 |
+When Gravitino creates a Lance dataset with blob columns, it picks the Lance
file format version for
+them: 2.2 for `lance.blob` and 2.1 for `lance.blob.legacy`. A table cannot
have both. Adding a blob
+column to an existing dataset fails with a clear error if the dataset's file
format version does not
+support it.
-For example, `External("lance.blob.v2(inline_size_threshold=4096,
dedicated_size_threshold=1048576)")`.
+Only the blob marker is recognized. Other blob metadata, such as
+`lance-encoding:blob-dedicated-size-threshold`, is not shown in the Gravitino
type, and columns created
+through Gravitino use Lance's defaults. A blob field in any other layout, for
example a blob v2 struct
+with the optional `position` and `size` fields or a legacy blob stored as
`Binary`, is returned as
+`External(arrow_field_json_str)` so that it keeps its exact definition.
-Legacy blob columns are rejected by Lance for file version 2.2 and later; use
blob v2 for new tables.
-Metadata other than the blob keys above is not kept. A blob field that does
not exactly match these
-layouts, for example a legacy blob stored as `Binary` or a threshold out of
the accepted range, is
-returned as `External(arrow_field_json_str)`.
+:::note
+Gravitino refreshes stored column types from the Lance dataset only when the
dataset changes. Until
+then, blob v2 columns loaded by an earlier Gravitino version still appear as a
plain `Struct`, and
+legacy blob columns stored as `External(arrow_field_json_str)` switch to
`External("lance.blob.legacy")`
+on the next refresh.
+:::
### Table Properties
diff --git
a/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceBlobTypes.java
b/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceBlobTypes.java
index 4a7835abc6..42d6a63814 100644
---
a/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceBlobTypes.java
+++
b/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceBlobTypes.java
@@ -19,11 +19,8 @@
package org.apache.gravitino.lance.common.ops.gravitino;
-import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
-import java.util.ArrayList;
-import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
@@ -32,118 +29,87 @@ import org.apache.arrow.vector.types.pojo.Field;
import org.apache.arrow.vector.types.pojo.FieldType;
/**
- * Converts Lance blob columns between their Arrow field form and a readable
catalog string used by
+ * Converts Lance blob columns between their Arrow field form and the catalog
strings used by
* Gravitino external types.
*
- * <p>Supported catalog strings:
- *
* <ul>
- * <li>{@code lance.blob.v1}: an Arrow {@code LargeBinary} field with
metadata {@code
- * lance-encoding:blob=true}.
- * <li>{@code lance.blob.v2(with_range=true, inline_size_threshold=N,
dedicated_size_threshold=N,
- * pack_file_size_threshold=N)}: an Arrow struct tagged with {@code
- * ARROW:extension:name=lance.blob.v2}. All parameters are optional;
without parameters the
- * parentheses are omitted.
+ * <li>{@code lance.blob}: a Lance blob v2 column, an Arrow {@code
Struct<data: LargeBinary, uri:
+ * Utf8>} tagged with {@code ARROW:extension:name=lance.blob.v2}. It
requires Lance file
+ * format version 2.2 or later.
+ * <li>{@code lance.blob.legacy}: a Lance legacy blob column, an Arrow
{@code LargeBinary} field
+ * with metadata {@code lance-encoding:blob=true}. Lance rejects it from
file format version
+ * 2.2 on.
* </ul>
*
- * <p>Only Lance blob metadata is recognized; other field metadata is not
represented. A blob field
- * that does not exactly match the canonical Lance layout is left to the Arrow
JSON representation.
+ * <p>Only the blob marker is recognized; other field metadata, such as
storage thresholds, is not
+ * represented. A blob field in any other layout is left to the Arrow JSON
representation so that it
+ * round-trips unchanged.
*/
final class LanceBlobTypes {
+ static final String BLOB = "lance.blob";
+ static final String LEGACY_BLOB = "lance.blob.legacy";
+
static final String BLOB_META_KEY = "lance-encoding:blob";
static final String ARROW_EXT_NAME_KEY = "ARROW:extension:name";
static final String BLOB_V2_EXT_NAME = "lance.blob.v2";
- static final String INLINE_SIZE_THRESHOLD_META_KEY =
"lance-encoding:blob-inline-size-threshold";
- static final String DEDICATED_SIZE_THRESHOLD_META_KEY =
- "lance-encoding:blob-dedicated-size-threshold";
- static final String PACK_FILE_SIZE_THRESHOLD_META_KEY =
- "lance-encoding:blob-pack-file-size-threshold";
-
- static final String V1 = "lance.blob.v1";
- static final String V2 = "lance.blob.v2";
-
- private static final String PREFIX = "lance.blob.";
- private static final String WITH_RANGE = "with_range";
- // Catalog string parameter name -> Arrow metadata key, in canonical output
order.
- private static final Map<String, String> THRESHOLD_PARAMS =
- ImmutableMap.of(
- "inline_size_threshold", INLINE_SIZE_THRESHOLD_META_KEY,
- "dedicated_size_threshold", DEDICATED_SIZE_THRESHOLD_META_KEY,
- "pack_file_size_threshold", PACK_FILE_SIZE_THRESHOLD_META_KEY);
+ /** The Lance file format version from which blob v2 is supported and legacy
blob is rejected. */
+ static final String BLOB_FILE_FORMAT_VERSION = "2.2";
- // Catalog string parameter name -> minimum accepted value, matching Lance's
validation.
- private static final Map<String, Long> THRESHOLD_MINIMUMS =
- ImmutableMap.of(
- "inline_size_threshold", 0L,
- "dedicated_size_threshold", 1L,
- "pack_file_size_threshold", 1L);
+ /** The Lance file format version used for datasets with legacy blob
columns. */
+ static final String LEGACY_BLOB_FILE_FORMAT_VERSION = "2.1";
- private static final ArrowType UINT64 = new ArrowType.Int(64, false);
-
- private static final List<Field> V2_MINIMAL_CHILDREN =
+ private static final List<Field> BLOB_CHILDREN =
ImmutableList.of(
- nullableChild("data", ArrowType.LargeBinary.INSTANCE),
- nullableChild("uri", ArrowType.Utf8.INSTANCE));
-
- private static final List<Field> V2_FULL_CHILDREN =
- ImmutableList.<Field>builder()
- .addAll(V2_MINIMAL_CHILDREN)
- .add(nullableChild("position", UINT64))
- .add(nullableChild("size", UINT64))
- .build();
-
- private static final String SUPPORTED_FORMATS =
- V1
- + ", "
- + V2
- + "("
- + WITH_RANGE
- + "=true, "
- + String.join("=N, ", THRESHOLD_PARAMS.keySet())
- + "=N)";
+ new Field("data", new FieldType(true,
ArrowType.LargeBinary.INSTANCE, null), null),
+ new Field("uri", new FieldType(true, ArrowType.Utf8.INSTANCE, null),
null));
private LanceBlobTypes() {}
/**
- * Returns whether the field carries Lance blob metadata, either legacy blob
or blob v2.
+ * Returns whether the field carries a Lance blob marker, either legacy blob
or blob v2. This
+ * matches Lance's own check, which only looks at whether the legacy key is
present.
*
* @param field The Arrow field.
* @return true if the field is a Lance blob field.
*/
static boolean isBlob(Field field) {
- Map<String, String> metadata = field.getMetadata();
- return metadata != null
- && (metadata.containsKey(BLOB_META_KEY)
- || BLOB_V2_EXT_NAME.equals(metadata.get(ARROW_EXT_NAME_KEY)));
+ return isBlobV2(field) || isLegacyBlob(field);
}
/**
- * Returns the catalog string of a canonical Lance blob field.
+ * Returns the catalog string of a Lance blob field in its standard layout.
*
* @param field The Arrow field.
- * @return The catalog string, or empty if the field is not a canonical
Lance blob.
+ * @return The catalog string, or empty if the field is not a Lance blob in
its standard layout.
*/
static Optional<String> toCatalogString(Field field) {
- Map<String, String> metadata = field.getMetadata();
- if (metadata == null || metadata.isEmpty()) {
+ if (field.getDictionary() != null) {
return Optional.empty();
}
- if (isCanonicalV1(field, metadata)) {
- return Optional.of(V1);
+ if (isBlobV2(field)
+ && field.getType() instanceof ArrowType.Struct
+ && childrenMatch(field.getChildren(), BLOB_CHILDREN)) {
+ return Optional.of(BLOB);
+ }
+ if (isLegacyBlob(field)
+ && field.getType() instanceof ArrowType.LargeBinary
+ && field.getChildren().isEmpty()
+ && "true".equals(field.getMetadata().get(BLOB_META_KEY))) {
+ return Optional.of(LEGACY_BLOB);
}
- return toV2CatalogString(field, metadata);
+ return Optional.empty();
}
/**
- * Returns whether the catalog string uses the Lance blob format.
+ * Returns whether the catalog string names a Lance blob type.
*
* @param catalogString The external type catalog string.
- * @return true if the string starts with {@code lance.blob.}.
+ * @return true if the string starts with {@code lance.blob}.
*/
static boolean isBlobCatalogString(String catalogString) {
- return catalogString != null && catalogString.trim().startsWith(PREFIX);
+ return catalogString != null && catalogString.trim().startsWith(BLOB);
}
/**
@@ -153,160 +119,122 @@ final class LanceBlobTypes {
* @param nullable Whether the field is nullable.
* @param catalogString The Lance blob catalog string.
* @return The Arrow field.
- * @throws IllegalArgumentException If the catalog string is not a valid
Lance blob type.
+ * @throws IllegalArgumentException If the catalog string is not a Lance
blob type.
*/
static Field toArrowField(String name, boolean nullable, String
catalogString) {
- String trimmed = catalogString.trim();
- if (trimmed.equals(V1)) {
- return new Field(
- name,
- new FieldType(
- nullable,
- ArrowType.LargeBinary.INSTANCE,
- null,
- ImmutableMap.of(BLOB_META_KEY, "true")),
- null);
- }
-
- Preconditions.checkArgument(
- trimmed.startsWith(V2),
- "Unsupported Lance blob type %s, expected one of: %s",
- trimmed,
- SUPPORTED_FORMATS);
- String rest = trimmed.substring(V2.length()).trim();
- Map<String, String> params = parseParams(rest, trimmed);
-
- boolean withRange = false;
- Map<String, String> metadata = new LinkedHashMap<>();
- metadata.put(ARROW_EXT_NAME_KEY, BLOB_V2_EXT_NAME);
- for (Map.Entry<String, String> param : params.entrySet()) {
- String key = param.getKey();
- String value = param.getValue();
- if (key.equals(WITH_RANGE)) {
- Preconditions.checkArgument(
- value.equals("true") || value.equals("false"),
- "Invalid value %s for %s in Lance blob type %s, expected true or
false",
- value,
- WITH_RANGE,
- trimmed);
- withRange = Boolean.parseBoolean(value);
- } else if (THRESHOLD_PARAMS.containsKey(key)) {
- metadata.put(THRESHOLD_PARAMS.get(key), parseThreshold(key, value,
trimmed));
- } else {
+ switch (catalogString.trim()) {
+ case BLOB:
+ return new Field(
+ name,
+ new FieldType(
+ nullable,
+ ArrowType.Struct.INSTANCE,
+ null,
+ ImmutableMap.of(ARROW_EXT_NAME_KEY, BLOB_V2_EXT_NAME)),
+ BLOB_CHILDREN);
+ case LEGACY_BLOB:
+ return new Field(
+ name,
+ new FieldType(
+ nullable,
+ ArrowType.LargeBinary.INSTANCE,
+ null,
+ ImmutableMap.of(BLOB_META_KEY, "true")),
+ null);
+ default:
throw new IllegalArgumentException(
String.format(
- "Unknown parameter %s in Lance blob type %s, expected one of:
%s",
- key, trimmed, SUPPORTED_FORMATS));
- }
+ "Unsupported Lance blob type %s, expected %s or %s",
+ catalogString.trim(), BLOB, LEGACY_BLOB));
}
-
- return new Field(
- name,
- new FieldType(nullable, ArrowType.Struct.INSTANCE, null, metadata),
- withRange ? V2_FULL_CHILDREN : V2_MINIMAL_CHILDREN);
}
- private static boolean isCanonicalV1(Field field, Map<String, String>
metadata) {
- return field.getType() instanceof ArrowType.LargeBinary
- && field.getDictionary() == null
- && field.getChildren().isEmpty()
- && "true".equals(metadata.get(BLOB_META_KEY));
- }
-
- private static Optional<String> toV2CatalogString(Field field, Map<String,
String> metadata) {
- if (!(field.getType() instanceof ArrowType.Struct)
- || field.getDictionary() != null
- || !BLOB_V2_EXT_NAME.equals(metadata.get(ARROW_EXT_NAME_KEY))) {
- return Optional.empty();
+ /**
+ * Returns the Lance file format version that a new dataset with these
fields must use.
+ *
+ * @param fields The top-level Arrow fields of the dataset.
+ * @return {@value #BLOB_FILE_FORMAT_VERSION} if any field contains a blob
v2 column, {@value
+ * #LEGACY_BLOB_FILE_FORMAT_VERSION} if any field contains a legacy blob
column, or empty if
+ * there is no blob column.
+ * @throws IllegalArgumentException If the fields contain both blob v2 and
legacy blob columns.
+ */
+ static Optional<String> requiredFileFormatVersion(List<Field> fields) {
+ boolean hasBlobV2 = fields.stream().anyMatch(field -> contains(field,
true));
+ boolean hasLegacyBlob = fields.stream().anyMatch(field -> contains(field,
false));
+ if (hasBlobV2 && hasLegacyBlob) {
+ throw new IllegalArgumentException(
+ String.format(
+ "A Lance table cannot mix %s and %s columns: %s requires file
format version %s or"
+ + " later, which does not support %s",
+ BLOB, LEGACY_BLOB, BLOB, BLOB_FILE_FORMAT_VERSION, LEGACY_BLOB));
}
-
- List<Field> children = field.getChildren();
- boolean withRange;
- if (childrenMatch(children, V2_MINIMAL_CHILDREN)) {
- withRange = false;
- } else if (childrenMatch(children, V2_FULL_CHILDREN)) {
- withRange = true;
- } else {
- return Optional.empty();
+ if (hasBlobV2) {
+ return Optional.of(BLOB_FILE_FORMAT_VERSION);
}
+ return hasLegacyBlob ? Optional.of(LEGACY_BLOB_FILE_FORMAT_VERSION) :
Optional.empty();
+ }
- List<String> params = new ArrayList<>();
- if (withRange) {
- params.add(WITH_RANGE + "=true");
+ /**
+ * Checks that a field can be added to a dataset with the given Lance file
format version.
+ *
+ * @param field The Arrow field to add.
+ * @param fileFormatVersion The dataset's Lance file format version, such as
{@code 2.1}.
+ * @throws IllegalArgumentException If the field contains a blob column the
version does not
+ * support.
+ */
+ static void checkFileFormatVersion(Field field, String fileFormatVersion) {
+ Optional<Boolean> supportsBlobV2 = supportsBlobV2(fileFormatVersion);
+ if (supportsBlobV2.isEmpty()) {
+ // Leave unknown versions to Lance.
+ return;
}
- for (Map.Entry<String, String> param : THRESHOLD_PARAMS.entrySet()) {
- String value = metadata.get(param.getValue());
- if (value == null) {
- continue;
- }
- if (!isCanonicalThreshold(param.getKey(), value)) {
- return Optional.empty();
- }
- params.add(param.getKey() + "=" + value);
+ if (!supportsBlobV2.get() && contains(field, true)) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Column %s of type %s requires Lance file format version %s or
later, but the"
+ + " dataset uses %s",
+ field.getName(), BLOB, BLOB_FILE_FORMAT_VERSION,
fileFormatVersion));
}
+ if (supportsBlobV2.get() && contains(field, false)) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Column %s of type %s is not supported by Lance file format
version %s or later,"
+ + " but the dataset uses %s; use %s instead",
+ field.getName(), LEGACY_BLOB, BLOB_FILE_FORMAT_VERSION,
fileFormatVersion, BLOB));
+ }
+ }
- return Optional.of(params.isEmpty() ? V2 : V2 + "(" + String.join(", ",
params) + ")");
+ private static boolean isBlobV2(Field field) {
+ Map<String, String> metadata = field.getMetadata();
+ return metadata != null &&
BLOB_V2_EXT_NAME.equals(metadata.get(ARROW_EXT_NAME_KEY));
}
- private static Map<String, String> parseParams(String rest, String
catalogString) {
- Map<String, String> params = new LinkedHashMap<>();
- if (rest.isEmpty()) {
- return params;
- }
- Preconditions.checkArgument(
- rest.startsWith("(") && rest.endsWith(")"),
- "Unsupported Lance blob type %s, expected one of: %s",
- catalogString,
- SUPPORTED_FORMATS);
- String body = rest.substring(1, rest.length() - 1).trim();
- if (body.isEmpty()) {
- return params;
- }
- for (String part : body.split(",", -1)) {
- String[] kv = part.split("=", -1);
- Preconditions.checkArgument(
- kv.length == 2 && !kv[0].trim().isEmpty() && !kv[1].trim().isEmpty(),
- "Invalid parameter '%s' in Lance blob type %s, expected key=value",
- part.trim(),
- catalogString);
- String key = kv[0].trim();
- Preconditions.checkArgument(
- params.put(key, kv[1].trim()) == null,
- "Duplicate parameter %s in Lance blob type %s",
- key,
- catalogString);
- }
- return params;
+ private static boolean isLegacyBlob(Field field) {
+ Map<String, String> metadata = field.getMetadata();
+ return metadata != null && metadata.containsKey(BLOB_META_KEY) &&
!isBlobV2(field);
}
- private static String parseThreshold(String key, String value, String
catalogString) {
- long threshold;
- try {
- threshold = Long.parseLong(value);
- } catch (NumberFormatException e) {
- throw new IllegalArgumentException(
- String.format(
- "Invalid value %s for %s in Lance blob type %s, expected an
integer",
- value, key, catalogString),
- e);
+ private static boolean contains(Field field, boolean blobV2) {
+ if (blobV2 ? isBlobV2(field) : isLegacyBlob(field)) {
+ return true;
}
- long min = THRESHOLD_MINIMUMS.get(key);
- Preconditions.checkArgument(
- threshold >= min,
- "Invalid value %s for %s in Lance blob type %s, expected an integer >=
%s",
- value,
- key,
- catalogString,
- min);
- return Long.toString(threshold);
+ return field.getChildren().stream().anyMatch(child -> contains(child,
blobV2));
}
- private static boolean isCanonicalThreshold(String key, String value) {
+ private static Optional<Boolean> supportsBlobV2(String fileFormatVersion) {
+ if (fileFormatVersion == null) {
+ return Optional.empty();
+ }
+ String[] parts = fileFormatVersion.trim().split("\\.");
+ if (parts.length != 2) {
+ return Optional.empty();
+ }
try {
- long threshold = Long.parseLong(value);
- return Long.toString(threshold).equals(value) && threshold >=
THRESHOLD_MINIMUMS.get(key);
+ int major = Integer.parseInt(parts[0]);
+ int minor = Integer.parseInt(parts[1]);
+ return Optional.of(major > 2 || (major == 2 && minor >= 2));
} catch (NumberFormatException e) {
- return false;
+ return Optional.empty();
}
}
@@ -329,8 +257,4 @@ final class LanceBlobTypes {
}
return true;
}
-
- private static Field nullableChild(String name, ArrowType type) {
- return new Field(name, new FieldType(true, type, null), null);
- }
}
diff --git
a/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceDataTypeConverter.java
b/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceDataTypeConverter.java
index a756f5acbc..1484baf929 100644
---
a/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceDataTypeConverter.java
+++
b/lance/lance-common/src/main/java/org/apache/gravitino/lance/common/ops/gravitino/LanceDataTypeConverter.java
@@ -47,6 +47,32 @@ public class LanceDataTypeConverter implements
DataTypeConverter<ArrowType, Fiel
public static final LanceDataTypeConverter CONVERTER = new
LanceDataTypeConverter();
private static final ObjectMapper mapper = new ObjectMapper();
+ /**
+ * Returns the Lance file format version that a new dataset with these
fields must use, so that
+ * its blob columns can be written: 2.2 for {@code lance.blob} and 2.1 for
{@code
+ * lance.blob.legacy}.
+ *
+ * @param fields The top-level Arrow fields of the dataset.
+ * @return The required file format version, or empty if the fields have no
blob column.
+ * @throws IllegalArgumentException If the fields mix {@code lance.blob} and
{@code
+ * lance.blob.legacy} columns.
+ */
+ public static Optional<String> requiredFileFormatVersion(List<Field> fields)
{
+ return LanceBlobTypes.requiredFileFormatVersion(fields);
+ }
+
+ /**
+ * Checks that a field can be added to a dataset with the given Lance file
format version.
+ *
+ * @param field The Arrow field to add.
+ * @param fileFormatVersion The dataset's Lance file format version, such as
{@code 2.1}.
+ * @throws IllegalArgumentException If the field contains a blob column the
version does not
+ * support.
+ */
+ public static void checkFileFormatVersion(Field field, String
fileFormatVersion) {
+ LanceBlobTypes.checkFileFormatVersion(field, fileFormatVersion);
+ }
+
public Field toArrowField(String name, Type type, boolean nullable) {
switch (type.name()) {
case LIST:
diff --git
a/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceBlobTypes.java
b/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceBlobTypes.java
index 46d81c90e7..e04e3bdd31 100644
---
a/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceBlobTypes.java
+++
b/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceBlobTypes.java
@@ -18,6 +18,7 @@
*/
package org.apache.gravitino.lance.common.ops.gravitino;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -40,56 +41,30 @@ import org.junit.jupiter.params.provider.ValueSource;
public class TestLanceBlobTypes {
+ private static final Map<String, String> LEGACY_META =
+ Map.of(LanceBlobTypes.BLOB_META_KEY, "true");
+
@Test
- void testV1RoundTrip() {
- Field field =
- new Field(
- "image",
- new FieldType(
- false,
- ArrowType.LargeBinary.INSTANCE,
- null,
- Map.of(LanceBlobTypes.BLOB_META_KEY, "true")),
- null);
+ void testBlobRoundTrip() {
+ Field field = blobField(standardChildren(), Map.of());
- assertEquals(Optional.of("lance.blob.v1"),
LanceBlobTypes.toCatalogString(field));
- assertEquals(field, LanceBlobTypes.toArrowField("image", false,
"lance.blob.v1"));
+ assertEquals(Optional.of("lance.blob"),
LanceBlobTypes.toCatalogString(field));
+ assertEquals(field, LanceBlobTypes.toArrowField("blob", true,
"lance.blob"));
}
@Test
- void testNonCanonicalV1FallsBack() {
- Map<String, String> blobMeta = Map.of(LanceBlobTypes.BLOB_META_KEY,
"true");
- assertEquals(
- Optional.empty(),
- LanceBlobTypes.toCatalogString(
- new Field("b", new FieldType(true, ArrowType.Binary.INSTANCE,
null, blobMeta), null)));
- assertEquals(
- Optional.empty(),
- LanceBlobTypes.toCatalogString(
- new Field(
- "b",
- new FieldType(
- true,
- ArrowType.LargeBinary.INSTANCE,
- new DictionaryEncoding(1L, false, new ArrowType.Int(32,
true)),
- blobMeta),
- null)));
- assertEquals(
- Optional.empty(),
- LanceBlobTypes.toCatalogString(
- new Field(
- "b",
- new FieldType(
- true,
- ArrowType.LargeBinary.INSTANCE,
- null,
- Map.of(LanceBlobTypes.BLOB_META_KEY, "yes")),
- null)));
+ void testLegacyBlobRoundTrip() {
+ Field field =
+ new Field(
+ "image", new FieldType(false, ArrowType.LargeBinary.INSTANCE,
null, LEGACY_META), null);
+
+ assertEquals(Optional.of("lance.blob.legacy"),
LanceBlobTypes.toCatalogString(field));
+ assertEquals(field, LanceBlobTypes.toArrowField("image", false,
"lance.blob.legacy"));
}
@Test
- void testNonBlobMetadataIsIgnored() {
- Field v1 =
+ void testNonLayoutMetadataIsIgnored() {
+ Field legacy =
new Field(
"b",
new FieldType(
@@ -98,21 +73,16 @@ public class TestLanceBlobTypes {
null,
Map.of(LanceBlobTypes.BLOB_META_KEY, "true", "other", "x")),
null);
- assertEquals(Optional.of("lance.blob.v1"),
LanceBlobTypes.toCatalogString(v1));
+ assertEquals(Optional.of("lance.blob.legacy"),
LanceBlobTypes.toCatalogString(legacy));
- // pyarrow extension types always export an ARROW:extension:metadata entry.
- Field v2 =
- v2Field(
- minimalChildren(),
+ // Storage thresholds and pyarrow's ARROW:extension:metadata entry.
+ Field blob =
+ blobField(
+ standardChildren(),
Map.of(
- "ARROW:extension:metadata",
- "",
- "lance-schema:unenforced-primary-key",
- "true",
- LanceBlobTypes.INLINE_SIZE_THRESHOLD_META_KEY,
- "16"));
- assertEquals(
- Optional.of("lance.blob.v2(inline_size_threshold=16)"),
LanceBlobTypes.toCatalogString(v2));
+ "ARROW:extension:metadata", "",
+ "lance-encoding:blob-dedicated-size-threshold", "16"));
+ assertEquals(Optional.of("lance.blob"),
LanceBlobTypes.toCatalogString(blob));
// Child metadata is not part of the blob v2 layout.
List<Field> childWithMetadata =
@@ -125,98 +95,69 @@ public class TestLanceBlobTypes {
null,
Map.of("lance-encoding:compression", "zstd")),
null),
- minimalChildren().get(1));
+ standardChildren().get(1));
assertEquals(
- Optional.of("lance.blob.v2"),
- LanceBlobTypes.toCatalogString(v2Field(childWithMetadata, Map.of())));
+ Optional.of("lance.blob"),
+ LanceBlobTypes.toCatalogString(blobField(childWithMetadata,
Map.of())));
}
@Test
- void testIsBlob() {
- assertTrue(LanceBlobTypes.isBlob(LanceBlobTypes.toArrowField("b", true,
"lance.blob.v1")));
- assertTrue(LanceBlobTypes.isBlob(LanceBlobTypes.toArrowField("b", true,
"lance.blob.v2")));
- assertTrue(
- LanceBlobTypes.isBlob(
- new Field(
- "b",
- new FieldType(
- true,
- ArrowType.Binary.INSTANCE,
- null,
- Map.of(LanceBlobTypes.BLOB_META_KEY, "true")),
- null)));
+ void testNonStandardLegacyBlobIsNotReadable() {
assertFalse(
- LanceBlobTypes.isBlob(
- new Field(
- "b",
- new FieldType(
- true,
- ArrowType.LargeBinary.INSTANCE,
- null,
- Map.of(LanceBlobTypes.INLINE_SIZE_THRESHOLD_META_KEY,
"1")),
- null)));
+ LanceBlobTypes.toCatalogString(
+ new Field(
+ "b", new FieldType(true, ArrowType.Binary.INSTANCE, null,
LEGACY_META), null))
+ .isPresent());
assertFalse(
- LanceBlobTypes.isBlob(
- new Field("b", new FieldType(true, ArrowType.LargeBinary.INSTANCE,
null), null)));
- }
-
- @Test
- void testV2MinimalRoundTrip() {
- Field field = v2Field(minimalChildren(), Map.of());
-
- assertEquals(Optional.of("lance.blob.v2"),
LanceBlobTypes.toCatalogString(field));
- assertEquals(field, LanceBlobTypes.toArrowField("blob", true,
"lance.blob.v2"));
- }
-
- @Test
- void testV2WithAllParamsRoundTrip() {
- Field field =
- v2Field(
- fullChildren(),
- Map.of(
- LanceBlobTypes.INLINE_SIZE_THRESHOLD_META_KEY, "0",
- LanceBlobTypes.DEDICATED_SIZE_THRESHOLD_META_KEY, "1048576",
- LanceBlobTypes.PACK_FILE_SIZE_THRESHOLD_META_KEY, "67108864"));
- String expected =
- "lance.blob.v2(with_range=true, inline_size_threshold=0, "
- + "dedicated_size_threshold=1048576,
pack_file_size_threshold=67108864)";
-
- assertEquals(Optional.of(expected), LanceBlobTypes.toCatalogString(field));
- assertEquals(field, LanceBlobTypes.toArrowField("blob", true, expected));
+ LanceBlobTypes.toCatalogString(
+ new Field(
+ "b",
+ new FieldType(
+ true,
+ ArrowType.LargeBinary.INSTANCE,
+ null,
+ Map.of(LanceBlobTypes.BLOB_META_KEY, "false")),
+ null))
+ .isPresent());
+ assertFalse(
+ LanceBlobTypes.toCatalogString(
+ new Field(
+ "b",
+ new FieldType(
+ true,
+ ArrowType.LargeBinary.INSTANCE,
+ new DictionaryEncoding(1L, false, new
ArrowType.Int(32, true)),
+ LEGACY_META),
+ null))
+ .isPresent());
}
@Test
- void testV2ParsesWithReorderedParamsAndWhitespace() {
- Field expected =
- v2Field(minimalChildren(),
Map.of(LanceBlobTypes.DEDICATED_SIZE_THRESHOLD_META_KEY, "10"));
-
- assertEquals(
- expected,
- LanceBlobTypes.toArrowField(
- "blob", true, " lance.blob.v2 ( dedicated_size_threshold = 10 ,
with_range=false ) "));
- assertEquals(
- v2Field(minimalChildren(), Map.of()),
- LanceBlobTypes.toArrowField("blob", true, "lance.blob.v2()"));
- }
+ void testNonStandardBlobIsNotReadable() {
+ // With the optional external range fields.
+ List<Field> fullChildren = new ArrayList<>(standardChildren());
+ fullChildren.add(
+ new Field("position", new FieldType(true, new ArrowType.Int(64,
false), null), null));
+ fullChildren.add(
+ new Field("size", new FieldType(true, new ArrowType.Int(64, false),
null), null));
+ assertFalse(LanceBlobTypes.toCatalogString(blobField(fullChildren,
Map.of())).isPresent());
- @Test
- void testNonCanonicalV2FallsBack() {
// Wrong child order.
- List<Field> reordered = new ArrayList<>(minimalChildren());
+ List<Field> reordered = new ArrayList<>(standardChildren());
Collections.reverse(reordered);
- assertFalse(LanceBlobTypes.toCatalogString(v2Field(reordered,
Map.of())).isPresent());
+ assertFalse(LanceBlobTypes.toCatalogString(blobField(reordered,
Map.of())).isPresent());
// Non-nullable child.
List<Field> nonNullable =
Arrays.asList(
new Field("data", new FieldType(false,
ArrowType.LargeBinary.INSTANCE, null), null),
- minimalChildren().get(1));
- assertFalse(LanceBlobTypes.toCatalogString(v2Field(nonNullable,
Map.of())).isPresent());
+ standardChildren().get(1));
+ assertFalse(LanceBlobTypes.toCatalogString(blobField(nonNullable,
Map.of())).isPresent());
// Dictionary-encoded child.
List<Field> dictionaryChild =
Arrays.asList(
- minimalChildren().get(0),
+ standardChildren().get(0),
new Field(
"uri",
new FieldType(
@@ -224,100 +165,128 @@ public class TestLanceBlobTypes {
ArrowType.Utf8.INSTANCE,
new DictionaryEncoding(1L, false, new ArrowType.Int(32,
true))),
null));
- assertFalse(LanceBlobTypes.toCatalogString(v2Field(dictionaryChild,
Map.of())).isPresent());
+ assertFalse(LanceBlobTypes.toCatalogString(blobField(dictionaryChild,
Map.of())).isPresent());
- // Dictionary-encoded struct.
- Field dictionaryEncoded =
- new Field(
- "blob",
- new FieldType(
- true,
- ArrowType.Struct.INSTANCE,
- new DictionaryEncoding(1L, false, new ArrowType.Int(32, true)),
- Map.of(LanceBlobTypes.ARROW_EXT_NAME_KEY,
LanceBlobTypes.BLOB_V2_EXT_NAME)),
- minimalChildren());
- assertFalse(LanceBlobTypes.toCatalogString(dictionaryEncoded).isPresent());
-
- // Thresholds out of the range accepted by the readable form.
- for (Map.Entry<String, String> threshold :
- Map.of(
- LanceBlobTypes.INLINE_SIZE_THRESHOLD_META_KEY, "-1",
- LanceBlobTypes.DEDICATED_SIZE_THRESHOLD_META_KEY, "0",
- LanceBlobTypes.PACK_FILE_SIZE_THRESHOLD_META_KEY, "0")
- .entrySet()) {
- assertFalse(
- LanceBlobTypes.toCatalogString(
- v2Field(minimalChildren(), Map.of(threshold.getKey(),
threshold.getValue())))
- .isPresent(),
- threshold.toString());
- }
-
- // Non-canonical threshold value.
+ // Not a struct.
assertFalse(
LanceBlobTypes.toCatalogString(
- v2Field(
- minimalChildren(),
- Map.of(LanceBlobTypes.INLINE_SIZE_THRESHOLD_META_KEY,
"04096")))
+ new Field(
+ "blob",
+ new FieldType(
+ true,
+ ArrowType.LargeBinary.INSTANCE,
+ null,
+ Map.of(LanceBlobTypes.ARROW_EXT_NAME_KEY,
LanceBlobTypes.BLOB_V2_EXT_NAME)),
+ null))
.isPresent());
}
@Test
- void testFieldWithoutBlobMetadataIsNotBlob() {
+ void testFieldWithoutBlobMarkerIsNotBlob() {
+ Field struct = new Field("s", new FieldType(true,
ArrowType.Struct.INSTANCE, null), null);
+ assertFalse(LanceBlobTypes.isBlob(struct));
+ assertFalse(LanceBlobTypes.toCatalogString(struct).isPresent());
assertFalse(
- LanceBlobTypes.toCatalogString(
- new Field("s", new FieldType(true, ArrowType.Struct.INSTANCE,
null), null))
- .isPresent());
+ LanceBlobTypes.isBlob(
+ new Field(
+ "b",
+ new FieldType(
+ true,
+ ArrowType.LargeBinary.INSTANCE,
+ null,
+ Map.of("lance-encoding:blob-dedicated-size-threshold",
"1")),
+ null)));
+ }
+
+ @Test
+ void testIsBlob() {
+ assertTrue(LanceBlobTypes.isBlob(LanceBlobTypes.toArrowField("b", true,
"lance.blob")));
+ assertTrue(LanceBlobTypes.isBlob(LanceBlobTypes.toArrowField("b", true,
"lance.blob.legacy")));
+ // Lance only checks that the legacy key is present.
+ assertTrue(
+ LanceBlobTypes.isBlob(
+ new Field(
+ "b",
+ new FieldType(
+ true,
+ ArrowType.Binary.INSTANCE,
+ null,
+ Map.of(LanceBlobTypes.BLOB_META_KEY, "false")),
+ null)));
}
@Test
void testIsBlobCatalogString() {
- assertTrue(LanceBlobTypes.isBlobCatalogString("lance.blob.v1"));
-
assertTrue(LanceBlobTypes.isBlobCatalogString("lance.blob.v2(with_range=true)"));
- assertTrue(LanceBlobTypes.isBlobCatalogString(" lance.blob.v1 "));
-
assertFalse(LanceBlobTypes.isBlobCatalogString("{\"name\":\"lance.blob.v1\"}"));
+ assertTrue(LanceBlobTypes.isBlobCatalogString("lance.blob"));
+ assertTrue(LanceBlobTypes.isBlobCatalogString(" lance.blob.legacy "));
+ assertTrue(LanceBlobTypes.isBlobCatalogString("lance.blob.v2"));
+
assertFalse(LanceBlobTypes.isBlobCatalogString("{\"name\":\"lance.blob\"}"));
assertFalse(LanceBlobTypes.isBlobCatalogString(null));
}
@ParameterizedTest
- @ValueSource(
- strings = {
- "lance.blob.v3",
- "lance.blob.v1(with_range=true)",
- "lance.blob.v2(unknown=1)",
- "lance.blob.v2(with_range=yes)",
- "lance.blob.v2(inline_size_threshold=abc)",
- "lance.blob.v2(inline_size_threshold=-1)",
- "lance.blob.v2(dedicated_size_threshold=0)",
- "lance.blob.v2(pack_file_size_threshold=0)",
- "lance.blob.v2(inline_size_threshold=1, inline_size_threshold=2)",
- "lance.blob.v2(inline_size_threshold)",
- "lance.blob.v2(inline_size_threshold=1",
- "lance.blob.v2x"
- })
+ @ValueSource(strings = {"lance.blob.v1", "lance.blob.v2", "lance.blob(x=1)",
"lance.blobs"})
void testInvalidCatalogStrings(String catalogString) {
assertThrows(
IllegalArgumentException.class,
() -> LanceBlobTypes.toArrowField("blob", true, catalogString));
}
- private static Field v2Field(List<Field> children, Map<String, String>
extraMetadata) {
+ @Test
+ void testRequiredFileFormatVersion() {
+ Field id = new Field("id", new FieldType(false, new ArrowType.Int(32,
true), null), null);
+ Field blob = LanceBlobTypes.toArrowField("blob", true, "lance.blob");
+ Field legacy = LanceBlobTypes.toArrowField("legacy", true,
"lance.blob.legacy");
+ Field nestedBlob =
+ new Field(
+ "record",
+ new FieldType(true, ArrowType.Struct.INSTANCE, null),
+ Collections.singletonList(blob));
+
+ assertEquals(Optional.empty(),
LanceBlobTypes.requiredFileFormatVersion(List.of(id)));
+ assertEquals(Optional.of("2.2"),
LanceBlobTypes.requiredFileFormatVersion(List.of(id, blob)));
+ assertEquals(
+ Optional.of("2.2"),
LanceBlobTypes.requiredFileFormatVersion(List.of(id, nestedBlob)));
+ assertEquals(Optional.of("2.1"),
LanceBlobTypes.requiredFileFormatVersion(List.of(id, legacy)));
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> LanceBlobTypes.requiredFileFormatVersion(List.of(blob, legacy)));
+ }
+
+ @Test
+ void testCheckFileFormatVersion() {
+ Field id = new Field("id", new FieldType(true, new ArrowType.Int(32,
true), null), null);
+ Field blob = LanceBlobTypes.toArrowField("blob", true, "lance.blob");
+ Field legacy = LanceBlobTypes.toArrowField("legacy", true,
"lance.blob.legacy");
+
+ assertDoesNotThrow(() -> LanceBlobTypes.checkFileFormatVersion(id, "2.1"));
+ assertDoesNotThrow(() -> LanceBlobTypes.checkFileFormatVersion(blob,
"2.2"));
+ assertDoesNotThrow(() -> LanceBlobTypes.checkFileFormatVersion(blob,
"2.3"));
+ assertDoesNotThrow(() -> LanceBlobTypes.checkFileFormatVersion(legacy,
"2.1"));
+ assertDoesNotThrow(() -> LanceBlobTypes.checkFileFormatVersion(legacy,
"0.1"));
+ // Unknown versions are left to Lance.
+ assertDoesNotThrow(() -> LanceBlobTypes.checkFileFormatVersion(blob,
"stable"));
+ assertDoesNotThrow(() -> LanceBlobTypes.checkFileFormatVersion(blob,
null));
+
+ IllegalArgumentException blobError =
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> LanceBlobTypes.checkFileFormatVersion(blob, "2.1"));
+ assertTrue(blobError.getMessage().contains("2.2"), blobError.getMessage());
+ assertThrows(
+ IllegalArgumentException.class, () ->
LanceBlobTypes.checkFileFormatVersion(legacy, "2.2"));
+ }
+
+ private static Field blobField(List<Field> children, Map<String, String>
extraMetadata) {
Map<String, String> metadata = new HashMap<>(extraMetadata);
metadata.put(LanceBlobTypes.ARROW_EXT_NAME_KEY,
LanceBlobTypes.BLOB_V2_EXT_NAME);
return new Field(
"blob", new FieldType(true, ArrowType.Struct.INSTANCE, null,
metadata), children);
}
- private static List<Field> minimalChildren() {
+ private static List<Field> standardChildren() {
return Arrays.asList(
new Field("data", new FieldType(true, ArrowType.LargeBinary.INSTANCE,
null), null),
new Field("uri", new FieldType(true, ArrowType.Utf8.INSTANCE, null),
null));
}
-
- private static List<Field> fullChildren() {
- List<Field> children = new ArrayList<>(minimalChildren());
- children.add(
- new Field("position", new FieldType(true, new ArrowType.Int(64,
false), null), null));
- children.add(new Field("size", new FieldType(true, new ArrowType.Int(64,
false), null), null));
- return children;
- }
}
diff --git
a/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceDataTypeConverter.java
b/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceDataTypeConverter.java
index 71558aa95e..478e3f9955 100644
---
a/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceDataTypeConverter.java
+++
b/lance/lance-common/src/test/java/org/apache/gravitino/lance/common/ops/gravitino/TestLanceDataTypeConverter.java
@@ -24,9 +24,10 @@ import static
org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
-import java.util.HashMap;
+import java.util.List;
import java.util.Map;
import java.util.function.Consumer;
import java.util.stream.Stream;
@@ -260,7 +261,8 @@ public class TestLanceDataTypeConverter {
}
@Test
- void testBlobV2ConvertsToReadableExternalTypeAndRoundTrips() {
+ void testBlobConvertsToReadableExternalTypeAndRoundTrips() {
+ // Storage thresholds and other non-layout metadata are not part of the
Gravitino type.
Field blobField =
new Field(
"blob",
@@ -270,27 +272,36 @@ public class TestLanceDataTypeConverter {
null,
Map.of(
"ARROW:extension:name", "lance.blob.v2",
- "lance-encoding:blob-inline-size-threshold", "4096",
"lance-encoding:blob-dedicated-size-threshold",
"1048576")),
Arrays.asList(
new Field("data", new FieldType(true,
ArrowType.LargeBinary.INSTANCE, null), null),
- new Field("uri", new FieldType(true, ArrowType.Utf8.INSTANCE,
null), null),
- new Field(
- "position", new FieldType(true, new ArrowType.Int(64,
false), null), null),
- new Field("size", new FieldType(true, new ArrowType.Int(64,
false), null), null)));
+ new Field("uri", new FieldType(true, ArrowType.Utf8.INSTANCE,
null), null)));
Type type = CONVERTER.toGravitino(blobField);
- assertInstanceOf(Types.ExternalType.class, type);
+ assertEquals(Types.ExternalType.of("lance.blob"), type);
assertEquals(
- "lance.blob.v2(with_range=true, inline_size_threshold=4096, "
- + "dedicated_size_threshold=1048576)",
- ((Types.ExternalType) type).catalogString());
- assertEquals(blobField, CONVERTER.toArrowField("blob", type, true));
+ LanceBlobTypes.toArrowField("blob", true, LanceBlobTypes.BLOB),
+ CONVERTER.toArrowField("blob", type, true));
+ }
+
+ @Test
+ void testLegacyBlobConvertsToReadableExternalTypeAndRoundTrips() {
+ Field legacyField =
+ new Field(
+ "image",
+ new FieldType(
+ false, ArrowType.LargeBinary.INSTANCE, null,
Map.of("lance-encoding:blob", "true")),
+ null);
+
+ Type type = CONVERTER.toGravitino(legacyField);
+
+ assertEquals(Types.ExternalType.of("lance.blob.legacy"), type);
+ assertEquals(legacyField, CONVERTER.toArrowField("image", type, false));
}
@Test
- void testNonCanonicalBlobConvertsToJsonExternalTypeAndRoundTrips() {
+ void testNonStandardBlobConvertsToJsonExternalTypeAndRoundTrips() {
// A legacy blob stored as Binary instead of LargeBinary.
assertJsonExternalTypeRoundTrips(
new Field(
@@ -299,16 +310,22 @@ public class TestLanceDataTypeConverter {
true, ArrowType.Binary.INSTANCE, null,
Map.of("lance-encoding:blob", "true")),
null));
- // A blob v2 struct whose threshold is out of range for the readable form.
- Field blobV2 =
- LanceBlobTypes.toArrowField("blob", true, LanceBlobTypes.V2 +
"(with_range=true)");
- Map<String, String> metadata = new HashMap<>(blobV2.getMetadata());
- metadata.put(LanceBlobTypes.DEDICATED_SIZE_THRESHOLD_META_KEY, "0");
+ // A legacy blob marker whose value is not "true"; Lance still treats it
as a blob.
assertJsonExternalTypeRoundTrips(
new Field(
- "blob",
- new FieldType(true, ArrowType.Struct.INSTANCE, null, metadata),
- blobV2.getChildren()));
+ "image",
+ new FieldType(
+ true, ArrowType.LargeBinary.INSTANCE, null,
Map.of("lance-encoding:blob", "false")),
+ null));
+
+ // A blob v2 struct with the optional external range fields.
+ Field blob = LanceBlobTypes.toArrowField("blob", true,
LanceBlobTypes.BLOB);
+ List<Field> fullChildren = new ArrayList<>(blob.getChildren());
+ fullChildren.add(
+ new Field("position", new FieldType(true, new ArrowType.Int(64,
false), null), null));
+ fullChildren.add(
+ new Field("size", new FieldType(true, new ArrowType.Int(64, false),
null), null));
+ assertJsonExternalTypeRoundTrips(new Field("blob", blob.getFieldType(),
fullChildren));
}
@Test
@@ -342,8 +359,7 @@ public class TestLanceDataTypeConverter {
@Test
void testStructWithBlobChildRoundTrips() {
- Field blobChild =
- LanceBlobTypes.toArrowField("image", true, LanceBlobTypes.V2 +
"(inline_size_threshold=1)");
+ Field blobChild = LanceBlobTypes.toArrowField("image", true,
LanceBlobTypes.BLOB);
Field structField =
new Field(
"record",
@@ -356,15 +372,13 @@ public class TestLanceDataTypeConverter {
Types.StructType structType = assertInstanceOf(Types.StructType.class,
type);
assertEquals(Types.LongType.get(), structType.fields()[0].type());
- assertEquals(
- Types.ExternalType.of("lance.blob.v2(inline_size_threshold=1)"),
- structType.fields()[1].type());
+ assertEquals(Types.ExternalType.of(LanceBlobTypes.BLOB),
structType.fields()[1].type());
assertEquals(structField, CONVERTER.toArrowField("record", type, true));
}
@Test
void testListWithBlobChildConvertsToNativeList() {
- String blobType = "lance.blob.v2(inline_size_threshold=16)";
+ String blobType = LanceBlobTypes.BLOB;
// Lance and pyarrow name list children "item"; Gravitino writes them back
as "element".
Field listField =
new Field(
@@ -416,7 +430,7 @@ public class TestLanceDataTypeConverter {
@Test
void testLargeAndFixedSizeListWithBlobChildRoundTrip() {
- Field blobChild = LanceBlobTypes.toArrowField("item", true,
LanceBlobTypes.V2);
+ Field blobChild = LanceBlobTypes.toArrowField("item", true,
LanceBlobTypes.BLOB);
assertJsonExternalTypeRoundTrips(
new Field(
"images",
@@ -440,7 +454,7 @@ public class TestLanceDataTypeConverter {
"images",
new FieldType(true, ArrowType.List.INSTANCE, null),
Collections.singletonList(
- LanceBlobTypes.toArrowField("element", true,
LanceBlobTypes.V2)))));
+ LanceBlobTypes.toArrowField("element", true,
LanceBlobTypes.BLOB)))));
Type type = CONVERTER.toGravitino(structField);
@@ -448,7 +462,7 @@ public class TestLanceDataTypeConverter {
Types.StructType.of(
Types.StructType.Field.of(
"images",
- Types.ListType.of(Types.ExternalType.of(LanceBlobTypes.V2),
true),
+ Types.ListType.of(Types.ExternalType.of(LanceBlobTypes.BLOB),
true),
true,
null)),
type);
@@ -471,12 +485,13 @@ public class TestLanceDataTypeConverter {
new FieldType(false, ArrowType.Utf8.INSTANCE,
null),
null),
LanceBlobTypes.toArrowField(
- MapVector.VALUE_NAME, true, LanceBlobTypes.V1)))));
+ MapVector.VALUE_NAME, true,
LanceBlobTypes.LEGACY_BLOB)))));
Type type = CONVERTER.toGravitino(mapField);
assertEquals(
- Types.MapType.of(Types.StringType.get(),
Types.ExternalType.of(LanceBlobTypes.V1), true),
+ Types.MapType.of(
+ Types.StringType.get(),
Types.ExternalType.of(LanceBlobTypes.LEGACY_BLOB), true),
type);
assertEquals(mapField, CONVERTER.toArrowField("images", type, true));
}
@@ -496,13 +511,14 @@ public class TestLanceDataTypeConverter {
}),
null),
Arrays.asList(
- LanceBlobTypes.toArrowField("image", true, LanceBlobTypes.V1),
+ LanceBlobTypes.toArrowField("image", true,
LanceBlobTypes.LEGACY_BLOB),
new Field("number", new FieldType(true, new ArrowType.Int(32,
true), null), null)));
Type type = CONVERTER.toGravitino(unionField);
assertEquals(
- Types.UnionType.of(Types.ExternalType.of(LanceBlobTypes.V1),
Types.IntegerType.get()),
+ Types.UnionType.of(
+ Types.ExternalType.of(LanceBlobTypes.LEGACY_BLOB),
Types.IntegerType.get()),
type);
Field blobChild = CONVERTER.toArrowField("value", type,
true).getChildren().get(0);
assertEquals(ArrowType.LargeBinary.INSTANCE, blobChild.getType());