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 7edfc468d53..b50b5c6a885 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 @@ -252,6 +252,10 @@ public Table createTable( Preconditions.checkArgument( StringUtils.isNotBlank(location), "Table location must be specified"); + // Validate every field before EXIST_OK or OVERWRITE can return, drop metadata, or purge a + // dataset. Conversion failures must never leave a partial mutation. + convertColumnsToArrowSchema(columns); + // Extract creation mode from properties CreationMode mode = Optional.ofNullable(properties.get(LANCE_CREATION_MODE)) 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 bf99a64814b..aa9e6f4ad2f 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 @@ -121,6 +121,36 @@ public void testCreateTableWithInvalidMode() { new Index[0])); } + @Test + public void testVariantTypeIsRejectedBeforeOverwriteMutation() { + NameIdentifier ident = NameIdentifier.of("catalog", "schema", "table"); + Column[] columns = {Column.of("payload", Types.VariantType.get(), "variant")}; + Map properties = + Map.of( + Table.PROPERTY_LOCATION, + tempDir.resolve("variant-overwrite").toString(), + LANCE_CREATION_MODE, + "OVERWRITE"); + + IllegalArgumentException exception = + Assertions.assertThrows( + IllegalArgumentException.class, + () -> + lanceTableOps.createTable( + ident, + columns, + null, + properties, + new Transform[0], + null, + new SortOrder[0], + new Index[0])); + + Assertions.assertTrue(exception.getMessage().contains("exact native representation")); + verify(lanceTableOps, never()).purgeTable(ident); + verify(lanceTableOps, never()).dropTable(ident); + } + @Test public void testLoadDeclaredTableSchemaFromLocation() throws Exception { NameIdentifier ident = NameIdentifier.of("schema", "table"); diff --git a/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/integration/test/CatalogGenericCatalogLanceIT.java b/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/integration/test/CatalogGenericCatalogLanceIT.java index f2dc40907ab..6c6505ac2f4 100644 --- a/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/integration/test/CatalogGenericCatalogLanceIT.java +++ b/catalogs/catalog-lakehouse-generic/src/test/java/org/apache/gravitino/catalog/lakehouse/lance/integration/test/CatalogGenericCatalogLanceIT.java @@ -44,6 +44,7 @@ import org.apache.arrow.vector.VarCharVector; import org.apache.arrow.vector.VectorSchemaRoot; import org.apache.arrow.vector.ipc.ArrowReader; +import org.apache.arrow.vector.types.TimeUnit; import org.apache.arrow.vector.types.pojo.ArrowType; import org.apache.arrow.vector.types.pojo.Field; import org.apache.commons.io.FileUtils; @@ -413,6 +414,203 @@ public void testCreateLanceTable() { RuntimeException.class, () -> catalog.asTableCatalog().loadTable(newNameIdentifier)); } + @Test + public void testNanosecondTimestampTypeRoundTrip() { + String phase2TableName = GravitinoITUtils.genRandomName("lance_timestamp_ns"); + NameIdentifier identifier = NameIdentifier.of(schemaName, phase2TableName); + String location = tempDirectory + "/" + phase2TableName; + Column[] columns = { + Column.of("timestamp_ns", Types.TimestampType.withoutTimeZone(9), "nanosecond timestamp"), + Column.of("timestamptz_ns", Types.TimestampType.withTimeZone(9), "nanosecond instant") + }; + Map properties = createProperties(); + properties.put(Table.PROPERTY_TABLE_FORMAT, LANCE_TABLE_FORMAT); + properties.put(Table.PROPERTY_LOCATION, location); + + Table created = + catalog + .asTableCatalog() + .createTable( + identifier, + columns, + "timestamp nanosecond round-trip", + properties, + Transforms.EMPTY_TRANSFORM, + null, + null); + Table loaded = catalog.asTableCatalog().loadTable(identifier); + + Assertions.assertEquals( + Types.TimestampType.withoutTimeZone(9), created.columns()[0].dataType()); + Assertions.assertEquals(Types.TimestampType.withTimeZone(9), created.columns()[1].dataType()); + Assertions.assertEquals(Types.TimestampType.withoutTimeZone(9), loaded.columns()[0].dataType()); + Assertions.assertEquals(Types.TimestampType.withTimeZone(9), loaded.columns()[1].dataType()); + + try (Dataset dataset = Dataset.open().uri(location).build()) { + List fields = dataset.getSchema().getFields(); + Assertions.assertEquals( + new ArrowType.Timestamp(TimeUnit.NANOSECOND, null), fields.get(0).getType()); + Assertions.assertEquals( + new ArrowType.Timestamp(TimeUnit.NANOSECOND, "UTC"), fields.get(1).getType()); + } + } + + @Test + public void testVariantOverwriteRejectedWithoutSideEffects() { + String phase2TableName = GravitinoITUtils.genRandomName("lance_variant_rejection"); + NameIdentifier identifier = NameIdentifier.of(schemaName, phase2TableName); + String location = tempDirectory + "/" + phase2TableName; + Map properties = createProperties(); + properties.put(Table.PROPERTY_TABLE_FORMAT, LANCE_TABLE_FORMAT); + properties.put(Table.PROPERTY_LOCATION, location); + properties.put(Table.PROPERTY_EXTERNAL, "true"); + Column[] originalColumns = { + Column.of("id", Types.IntegerType.get(), "original integer column") + }; + + catalog + .asTableCatalog() + .createTable( + identifier, + originalColumns, + "original table", + properties, + Transforms.EMPTY_TRANSFORM, + null, + null); + + Map overwriteProperties = Maps.newHashMap(properties); + overwriteProperties.put(LANCE_CREATION_MODE, "OVERWRITE"); + IllegalArgumentException exception = + Assertions.assertThrows( + IllegalArgumentException.class, + () -> + catalog + .asTableCatalog() + .createTable( + identifier, + new Column[] {Column.of("payload", Types.VariantType.get(), "variant")}, + "invalid overwrite", + overwriteProperties, + Transforms.EMPTY_TRANSFORM, + null, + null)); + + Assertions.assertTrue(exception.getMessage().contains("exact native representation")); + Assertions.assertTrue(catalog.asTableCatalog().tableExists(identifier)); + Assertions.assertEquals( + Types.IntegerType.get(), + catalog.asTableCatalog().loadTable(identifier).columns()[0].dataType()); + try (Dataset dataset = Dataset.open().uri(location).build()) { + Assertions.assertEquals( + new ArrowType.Int(32, true), dataset.getSchema().getFields().get(0).getType()); + } + } + + @Test + public void testUnknownTypeRoundTrip() { + String phase2TableName = GravitinoITUtils.genRandomName("lance_unknown"); + NameIdentifier identifier = NameIdentifier.of(schemaName, phase2TableName); + String location = tempDirectory + "/" + phase2TableName; + Column[] columns = {Column.of("unknown_value", Types.NullType.get(), "null-only value")}; + Map properties = createProperties(); + properties.put(Table.PROPERTY_TABLE_FORMAT, LANCE_TABLE_FORMAT); + properties.put(Table.PROPERTY_LOCATION, location); + + Table created = + catalog + .asTableCatalog() + .createTable( + identifier, + columns, + "unknown type round-trip", + properties, + Transforms.EMPTY_TRANSFORM, + null, + null); + Table loaded = catalog.asTableCatalog().loadTable(identifier); + + Assertions.assertEquals(Types.NullType.get(), created.columns()[0].dataType()); + Assertions.assertEquals(Types.NullType.get(), loaded.columns()[0].dataType()); + Assertions.assertTrue(created.columns()[0].nullable()); + try (Dataset dataset = Dataset.open().uri(location).build()) { + Field field = dataset.getSchema().getFields().get(0); + Assertions.assertEquals(ArrowType.Null.INSTANCE, field.getType()); + Assertions.assertTrue(field.isNullable()); + } + } + + @Test + public void testGeometryTypeRoundTrip() { + String phase2TableName = GravitinoITUtils.genRandomName("lance_geometry"); + NameIdentifier identifier = NameIdentifier.of(schemaName, phase2TableName); + String location = tempDirectory + "/" + phase2TableName; + Types.GeometryType geometry = Types.GeometryType.of("EPSG:3857"); + Column[] columns = {Column.of("shape", geometry, "planar WKB geometry")}; + Map properties = createProperties(); + properties.put(Table.PROPERTY_TABLE_FORMAT, LANCE_TABLE_FORMAT); + properties.put(Table.PROPERTY_LOCATION, location); + + Table created = + catalog + .asTableCatalog() + .createTable( + identifier, + columns, + "geometry type round-trip", + properties, + Transforms.EMPTY_TRANSFORM, + null, + null); + Table loaded = catalog.asTableCatalog().loadTable(identifier); + + Assertions.assertEquals(geometry, created.columns()[0].dataType()); + Assertions.assertEquals(geometry, loaded.columns()[0].dataType()); + try (Dataset dataset = Dataset.open().uri(location).build()) { + Field field = dataset.getSchema().getFields().get(0); + Assertions.assertEquals(ArrowType.Binary.INSTANCE, field.getType()); + Assertions.assertEquals("geoarrow.wkb", field.getMetadata().get("ARROW:extension:name")); + Assertions.assertEquals( + "{\"crs\":\"EPSG:3857\"}", field.getMetadata().get("ARROW:extension:metadata")); + } + } + + @Test + public void testGeographyTypeRoundTrip() { + String phase2TableName = GravitinoITUtils.genRandomName("lance_geography"); + NameIdentifier identifier = NameIdentifier.of(schemaName, phase2TableName); + String location = tempDirectory + "/" + phase2TableName; + Types.GeographyType geography = Types.GeographyType.of("EPSG:4326", "karney"); + Column[] columns = {Column.of("shape", geography, "ellipsoidal WKB geography")}; + Map properties = createProperties(); + properties.put(Table.PROPERTY_TABLE_FORMAT, LANCE_TABLE_FORMAT); + properties.put(Table.PROPERTY_LOCATION, location); + + Table created = + catalog + .asTableCatalog() + .createTable( + identifier, + columns, + "geography type round-trip", + properties, + Transforms.EMPTY_TRANSFORM, + null, + null); + Table loaded = catalog.asTableCatalog().loadTable(identifier); + + Assertions.assertEquals(geography, created.columns()[0].dataType()); + Assertions.assertEquals(geography, loaded.columns()[0].dataType()); + try (Dataset dataset = Dataset.open().uri(location).build()) { + Field field = dataset.getSchema().getFields().get(0); + Assertions.assertEquals(ArrowType.Binary.INSTANCE, field.getType()); + Assertions.assertEquals("geoarrow.wkb", field.getMetadata().get("ARROW:extension:name")); + Assertions.assertEquals( + "{\"crs\":\"EPSG:4326\",\"edges\":\"karney\"}", + field.getMetadata().get("ARROW:extension:metadata")); + } + } + @Test void testLanceTableFormat() { String tableName = GravitinoITUtils.genRandomName(TABLE_PREFIX); diff --git a/docs/lakehouse-generic-lance-table.md b/docs/lakehouse-generic-lance-table.md index cdbd0bb79ad..d241c037b1e 100644 --- a/docs/lakehouse-generic-lance-table.md +++ b/docs/lakehouse-generic-lance-table.md @@ -70,12 +70,37 @@ Lance uses Apache Arrow for table schemas. The following table shows type mappin | `Timestamp_tz(3)` | `TimestampType Millisecond withUtc` | | `Timestamp_tz(9)` | `TimestampType Nanosecond withUtc` | | `Time`/`Time(9)` | `Time Nanosecond` | -| `Null` | `Null` | +| `Unknown` (`NullType`) | `Null` (nullable columns only) | +| `Variant` | Rejected before mutation | +| `Geometry(crs)` | `geoarrow.wkb` with planar CRS metadata | +| `Geography(crs, algorithm)` | `geoarrow.wkb` with CRS and edge metadata | | `Fixed(n)` | `Fixed-Size Binary(n)` | | `Interval_year` | Not supported by Lance | | `Interval_day` | `Duration(Microsecond)` | | `External(arrow_field_json_str)` | Any Arrow Field | +`Timestamp(9)` and `Timestamp_tz(9)` round-trip losslessly through Lance as nanosecond Arrow +timestamps. Gravitino `Timestamp_tz` has time-zone-aware instant semantics but does not carry a +zone identifier, so the connector writes the canonical Arrow `UTC` identifier. Native Arrow +timestamps with another zone identifier are preserved as `External` instead of silently rewriting +that identifier. + +Lance 6.0 with Arrow 18 has no exact native representation for Gravitino `Variant`, so table +creation rejects it before creating, replacing, or deleting table data. Arrow extension fields, +including newer `arrow.parquet.variant` fields, remain lossless `External` types. + +Gravitino `Unknown` is a nullable, null-only placeholder whose concrete type may be assigned during +schema evolution. It round-trips exactly as Arrow `Null`; a non-nullable Unknown column is rejected +because it cannot contain a valid value. + +Gravitino `Geometry` uses WKB with planar edges and CRS type metadata. Lance preserves the same +semantics as a GeoArrow 0.2 `geoarrow.wkb` extension field backed by Arrow `Binary`; both textual +CRS identifiers and PROJJSON round-trip without losing the CRS. + +Gravitino `Geography` uses the same GeoArrow WKB storage and records its spherical or spheroidal +edge algorithm in the GeoArrow `edges` metadata. All five Gravitino algorithms (`spherical`, +`vincenty`, `thomas`, `andoyer`, and `karney`) round-trip with the CRS. + ### External Types For Arrow types not natively mapped in Gravitino, use the `External(arrow_field_json_str)` type, which accepts a JSON string representation of an Arrow `Field`. 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 52d52d38fbd..9c8843e4eee 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 @@ -19,11 +19,16 @@ package org.apache.gravitino.lance.common.ops.gravitino; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.base.Preconditions; import com.google.common.collect.Lists; import java.util.Arrays; import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; import org.apache.arrow.vector.complex.MapVector; import org.apache.arrow.vector.types.DateUnit; import org.apache.arrow.vector.types.FloatingPointPrecision; @@ -44,6 +49,10 @@ public class LanceDataTypeConverter implements DataTypeConverter { public static final LanceDataTypeConverter CONVERTER = new LanceDataTypeConverter(); + private static final String ARROW_EXTENSION_NAME = "ARROW:extension:name"; + private static final String ARROW_EXTENSION_METADATA = "ARROW:extension:metadata"; + private static final String GEOARROW_WKB_EXTENSION_NAME = "geoarrow.wkb"; + private static final Set GEOARROW_METADATA_KEYS = Set.of("crs", "edges"); private static final ObjectMapper mapper = new ObjectMapper(); public Field toArrowField(String name, Type type, boolean nullable) { @@ -132,6 +141,16 @@ public Field toArrowField(String name, Type type, boolean nullable) { field.isNullable()); return field; + case NULL: + Preconditions.checkArgument(nullable, "Lance Arrow Null columns must be nullable"); + return Field.nullable(name, ArrowType.Null.INSTANCE); + + case GEOMETRY: + return toGeoArrowField(name, (Types.GeometryType) type, nullable); + + case GEOGRAPHY: + return toGeoArrowField(name, (Types.GeographyType) type, nullable); + default: // non-complex type FieldType fieldType = new FieldType(nullable, fromGravitino(type), null); @@ -196,6 +215,15 @@ public ArrowType fromGravitino(Type type) { case FIXED: FixedType fixedType = (FixedType) type; return new ArrowType.FixedSizeBinary(fixedType.length()); + case VARIANT: + throw new IllegalArgumentException( + "Lance 6.0 and Arrow 18 do not provide an exact native representation for Gravitino Variant"); + case GEOMETRY: + throw new IllegalArgumentException( + "Gravitino Geometry requires GeoArrow field metadata; use toArrowField for Lance schemas"); + case GEOGRAPHY: + throw new IllegalArgumentException( + "Gravitino Geography requires GeoArrow field metadata; use toArrowField for Lance schemas"); default: throw new UnsupportedOperationException("Unsupported Gravitino type: " + type.name()); } @@ -204,6 +232,11 @@ public ArrowType fromGravitino(Type type) { @Override public Type toGravitino(Field arrowField) { FieldType fieldType = arrowField.getFieldType(); + if (fieldType.getMetadata() != null + && fieldType.getMetadata().containsKey(ARROW_EXTENSION_NAME)) { + return toKnownExtensionType(arrowField).orElseGet(() -> toExternalType(arrowField)); + } + switch (fieldType.getType().getTypeID()) { case Map: Field structField = arrowField.getChildren().get(0); @@ -284,10 +317,13 @@ public Type toGravitino(Field arrowField) { case MICROSECOND -> 6; case NANOSECOND -> 9; }; - boolean hasTimeZone = timestampType.getTimezone() != null; - return hasTimeZone - ? Types.TimestampType.withTimeZone(precision) - : Types.TimestampType.withoutTimeZone(precision); + if (timestampType.getTimezone() == null) { + return Types.TimestampType.withoutTimeZone(precision); + } + if ("UTC".equals(timestampType.getTimezone())) { + return Types.TimestampType.withTimeZone(precision); + } + break; case Time: ArrowType.Time timeType = (ArrowType.Time) fieldType.getType(); if (timeType.getUnit() == TimeUnit.NANOSECOND && timeType.getBitWidth() == 8 * 8) { @@ -315,12 +351,105 @@ public Type toGravitino(Field arrowField) { // fallthrough } - String typeString; + return toExternalType(arrowField); + } + + private Field toGeoArrowField(String name, Types.GeometryType type, boolean nullable) { + return toGeoArrowField(name, type.crs(), Optional.empty(), nullable); + } + + private Field toGeoArrowField(String name, Types.GeographyType type, boolean nullable) { + return toGeoArrowField(name, type.crs(), Optional.of(type.algorithm()), nullable); + } + + private Field toGeoArrowField( + String name, String crs, Optional edgeAlgorithm, boolean nullable) { + ObjectNode extensionMetadata = mapper.createObjectNode(); + putCrs(extensionMetadata, crs); + edgeAlgorithm.ifPresent(algorithm -> extensionMetadata.put("edges", algorithm)); + FieldType fieldType = + new FieldType( + nullable, + ArrowType.Binary.INSTANCE, + null, + Map.of( + ARROW_EXTENSION_NAME, + GEOARROW_WKB_EXTENSION_NAME, + ARROW_EXTENSION_METADATA, + extensionMetadata.toString())); + return new Field(name, fieldType, null); + } + + private Optional toKnownExtensionType(Field arrowField) { + Map metadata = arrowField.getFieldType().getMetadata(); + if (!GEOARROW_WKB_EXTENSION_NAME.equals(metadata.get(ARROW_EXTENSION_NAME)) + || metadata.size() != 2 + || !metadata.containsKey(ARROW_EXTENSION_METADATA) + || !(arrowField.getType() instanceof ArrowType.Binary) + || arrowField.getFieldType().getDictionary() != null + || !arrowField.getChildren().isEmpty()) { + return Optional.empty(); + } + + String extensionMetadata = metadata.get(ARROW_EXTENSION_METADATA); + if (extensionMetadata == null) { + return Optional.empty(); + } + try { + JsonNode metadataNode = mapper.readTree(extensionMetadata); + if (!metadataNode.isObject() || !hasOnlyGeoArrowMetadataKeys(metadataNode)) { + return Optional.empty(); + } + String crs = readCrs(metadataNode.get("crs")); + if (crs == null) { + return Optional.empty(); + } + JsonNode edgesNode = metadataNode.get("edges"); + if (edgesNode == null) { + return Optional.of(Types.GeometryType.of(crs)); + } + if (!edgesNode.isTextual()) { + return Optional.empty(); + } + return Optional.of(Types.GeographyType.of(crs, edgesNode.asText())); + } catch (Exception e) { + return Optional.empty(); + } + } + + private boolean hasOnlyGeoArrowMetadataKeys(JsonNode metadataNode) { + return Lists.newArrayList(metadataNode.fieldNames()).stream() + .allMatch(GEOARROW_METADATA_KEYS::contains); + } + + private void putCrs(ObjectNode metadataNode, String crs) { + try { + JsonNode parsedCrs = mapper.readTree(crs); + if (parsedCrs.isObject()) { + metadataNode.set("crs", parsedCrs); + return; + } + } catch (Exception ignored) { + // A CRS string is a valid GeoArrow fallback when it is not PROJJSON. + } + metadataNode.put("crs", crs); + } + + private String readCrs(JsonNode crsNode) { + if (crsNode == null) { + return null; + } + if (crsNode.isTextual() && !crsNode.asText().isEmpty()) { + return crsNode.asText(); + } + return crsNode.isObject() ? crsNode.toString() : null; + } + + private Types.ExternalType toExternalType(Field arrowField) { try { - typeString = mapper.writeValueAsString(arrowField); + return Types.ExternalType.of(mapper.writeValueAsString(arrowField)); } catch (Exception e) { throw new RuntimeException("Failed to serialize Arrow field to string.", e); } - return Types.ExternalType.of(typeString); } } 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 9908f8feffd..04c497b10c0 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 @@ -26,6 +26,8 @@ import java.util.Arrays; import java.util.Collections; +import java.util.List; +import java.util.Map; import java.util.function.Consumer; import java.util.stream.Stream; import org.apache.arrow.vector.complex.MapVector; @@ -207,6 +209,195 @@ public void testFromGravitinoTimestampWithTz() { assertEquals("UTC", tsArrow.getTimezone()); } + @Test + public void testNanosecondTimestampRoundTrip() { + Types.TimestampType timestampNs = Types.TimestampType.withoutTimeZone(9); + Types.TimestampType timestampTzNs = Types.TimestampType.withTimeZone(9); + + Field timestampField = CONVERTER.toArrowField("timestamp_ns", timestampNs, true); + Field timestampTzField = CONVERTER.toArrowField("timestamptz_ns", timestampTzNs, true); + + assertEquals( + new ArrowType.Timestamp(TimeUnit.NANOSECOND, null), + timestampField.getFieldType().getType()); + assertEquals( + new ArrowType.Timestamp(TimeUnit.NANOSECOND, "UTC"), + timestampTzField.getFieldType().getType()); + assertEquals(timestampNs, CONVERTER.toGravitino(timestampField)); + assertEquals(timestampTzNs, CONVERTER.toGravitino(timestampTzField)); + } + + @Test + public void testNonUtcTimestampPreservesTimezoneAsExternalType() { + Field field = + Field.nullable( + "event_time", new ArrowType.Timestamp(TimeUnit.NANOSECOND, "America/Los_Angeles")); + + Type converted = CONVERTER.toGravitino(field); + + assertInstanceOf(Types.ExternalType.class, converted); + assertEquals(field, CONVERTER.toArrowField("event_time", converted, true)); + } + + @Test + public void testVariantTypeIsRejectedAndArrowExtensionsRemainExternal() { + IllegalArgumentException exception = + assertThrows( + IllegalArgumentException.class, () -> CONVERTER.fromGravitino(Types.VariantType.get())); + assertTrue(exception.getMessage().contains("exact native representation")); + + Field nativeVariant = + new Field( + "payload", + new FieldType( + true, + ArrowType.Struct.INSTANCE, + null, + Map.of("ARROW:extension:name", "arrow.parquet.variant")), + List.of( + Field.notNullable("metadata", ArrowType.Binary.INSTANCE), + Field.nullable("value", ArrowType.Binary.INSTANCE))); + Type converted = CONVERTER.toGravitino(nativeVariant); + + assertInstanceOf(Types.ExternalType.class, converted); + assertEquals(nativeVariant, CONVERTER.toArrowField("payload", converted, true)); + } + + @Test + public void testUnknownTypeRoundTripAndNullability() { + Field field = CONVERTER.toArrowField("unknown_value", Types.NullType.get(), true); + + assertEquals(ArrowType.Null.INSTANCE, field.getType()); + assertTrue(field.isNullable()); + assertEquals(Types.NullType.get(), CONVERTER.toGravitino(field)); + + IllegalArgumentException exception = + assertThrows( + IllegalArgumentException.class, + () -> CONVERTER.toArrowField("unknown_value", Types.NullType.get(), false)); + assertTrue(exception.getMessage().contains("Null columns must be nullable")); + } + + @Test + public void testGeometryTypeRoundTripAsGeoArrowWkb() { + Types.GeometryType geometry = Types.GeometryType.of("EPSG:3857"); + + Field field = CONVERTER.toArrowField("shape", geometry, true); + + assertEquals(ArrowType.Binary.INSTANCE, field.getType()); + assertEquals("geoarrow.wkb", field.getMetadata().get("ARROW:extension:name")); + assertEquals("{\"crs\":\"EPSG:3857\"}", field.getMetadata().get("ARROW:extension:metadata")); + assertEquals(geometry, CONVERTER.toGravitino(field)); + + IllegalArgumentException exception = + assertThrows(IllegalArgumentException.class, () -> CONVERTER.fromGravitino(geometry)); + assertTrue(exception.getMessage().contains("requires GeoArrow field metadata")); + } + + @Test + public void testGeometryProjJsonCrsRoundTrip() { + String projJson = "{\"type\":\"ProjectedCRS\",\"name\":\"WGS 84 / Pseudo-Mercator\"}"; + Types.GeometryType geometry = Types.GeometryType.of(projJson); + + Field field = CONVERTER.toArrowField("shape", geometry, true); + + assertEquals("{\"crs\":" + projJson + "}", field.getMetadata().get("ARROW:extension:metadata")); + assertEquals(geometry, CONVERTER.toGravitino(field)); + } + + @ParameterizedTest + @CsvSource({"spherical", "vincenty", "thomas", "andoyer", "karney"}) + public void testGeographyTypeRoundTripAsGeoArrowWkb(String edgeAlgorithm) { + Types.GeographyType geography = Types.GeographyType.of("EPSG:4326", edgeAlgorithm); + + Field field = CONVERTER.toArrowField("shape", geography, true); + + assertEquals(ArrowType.Binary.INSTANCE, field.getType()); + assertEquals("geoarrow.wkb", field.getMetadata().get("ARROW:extension:name")); + assertEquals( + String.format("{\"crs\":\"EPSG:4326\",\"edges\":\"%s\"}", edgeAlgorithm), + field.getMetadata().get("ARROW:extension:metadata")); + assertEquals(geography, CONVERTER.toGravitino(field)); + } + + @Test + public void testUnsupportedGeoArrowEdgesRemainExternal() { + Field field = + new Field( + "shape", + new FieldType( + true, + ArrowType.Binary.INSTANCE, + null, + Map.of( + "ARROW:extension:name", + "geoarrow.wkb", + "ARROW:extension:metadata", + "{\"crs\":\"EPSG:4326\",\"edges\":\"rhumb\"}")), + null); + + Type converted = CONVERTER.toGravitino(field); + + assertInstanceOf(Types.ExternalType.class, converted); + assertEquals(field, CONVERTER.toArrowField("shape", converted, true)); + } + + @Test + public void testGeoArrowRepresentationMetadataRemainsExternal() { + Field field = + new Field( + "shape", + new FieldType( + true, + ArrowType.Binary.INSTANCE, + null, + Map.of( + "ARROW:extension:name", + "geoarrow.wkb", + "ARROW:extension:metadata", + "{\"crs\":\"EPSG:4326\",\"crs_type\":\"authority_code\",\"edges\":\"karney\"}")), + null); + + Type converted = CONVERTER.toGravitino(field); + + assertInstanceOf(Types.ExternalType.class, converted); + assertEquals(field, CONVERTER.toArrowField("shape", converted, true)); + } + + @Test + public void testGeoArrowCustomFieldMetadataRemainsExternal() { + Field field = + new Field( + "shape", + new FieldType( + true, + ArrowType.Binary.INSTANCE, + null, + Map.of( + "ARROW:extension:name", + "geoarrow.wkb", + "ARROW:extension:metadata", + "{\"crs\":\"EPSG:4326\",\"edges\":\"karney\"}", + "producer", + "native-client")), + null); + + Type converted = CONVERTER.toGravitino(field); + + assertInstanceOf(Types.ExternalType.class, converted); + assertEquals(field, CONVERTER.toArrowField("shape", converted, true)); + } + + @Test + public void testGeographyRequiresFieldMetadata() { + IllegalArgumentException exception = + assertThrows( + IllegalArgumentException.class, + () -> CONVERTER.fromGravitino(Types.GeographyType.crs84())); + + assertTrue(exception.getMessage().contains("requires GeoArrow field metadata")); + } + @Test public void testExternalTypeConversion() { String expectedColumnName = "col_name";