This is an automated email from the ASF dual-hosted git repository.
mchades pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new 816250fc2b [#12600] feat(store): Add Semantic Model create and load
persistence (#12602)
816250fc2b is described below
commit 816250fc2bbf92030a9d93b470e33fe68284ca64
Author: mchades <[email protected]>
AuthorDate: Tue Sep 22 19:45:19 2026 +0800
[#12600] feat(store): Add Semantic Model create and load persistence
(#12602)
### What changes were proposed in this pull request?
Add relational persistence for creating and loading Semantic Models,
including persistence objects, definition SerDe, mappers, SQL providers,
service routing, JDBC backend integration, and cache registration.
The feature delta is 20 files, with 1,727 insertions and 5 deletions.
### Why are the changes needed?
Semantic Models need durable relational persistence with correct
definition, property, and transaction semantics.
Fix: #12600
### Does this PR introduce _any_ user-facing change?
No new public API, REST, OpenAPI, or client surface is introduced.
### How was this patch tested?
- Spotless, Javadoc, compilation, test compilation, and `git diff
--check` passed.
- `TestSemanticModelJDBCBackend` passed 9/9 across H2, MySQL, and
PostgreSQL.
- `TestBaseEntityCache` passed 6/6.
---
.../apache/gravitino/cache/BaseEntityCache.java | 4 +-
.../gravitino/cache/SupportsEntityStoreCache.java | 6 +-
.../gravitino/storage/relational/JDBCBackend.java | 7 +
.../RelationalEntityStoreIdResolver.java | 15 +-
.../relational/mapper/SemanticModelMetaMapper.java | 143 ++++++
.../SemanticModelMetaSQLProviderFactory.java | 107 ++++
.../mapper/SemanticModelVersionInfoMapper.java | 37 ++
...SemanticModelVersionInfoSQLProviderFactory.java | 60 +++
.../provider/DefaultMapperPackageProvider.java | 4 +
.../base/SemanticModelMetaBaseSQLProvider.java | 208 ++++++++
.../SemanticModelVersionInfoBaseSQLProvider.java | 45 ++
.../SemanticModelMetaPostgreSQLProvider.java | 24 +
...SemanticModelVersionInfoPostgreSQLProvider.java | 25 +
.../storage/relational/po/SemanticModelPO.java | 228 +++++++++
.../relational/po/SemanticModelVersionInfoPO.java | 91 ++++
.../service/SemanticModelMetaService.java | 242 +++++++++
.../service/SemanticModelPOStorageOps.java | 87 ++++
.../relational/TestSemanticModelJDBCBackend.java | 568 +++++++++++++++++++++
.../relational/service/TestSchemaMetaService.java | 34 ++
.../gravitino-entity-cache-multinode-design.md | 14 +-
20 files changed, 1937 insertions(+), 12 deletions(-)
diff --git a/core/src/main/java/org/apache/gravitino/cache/BaseEntityCache.java
b/core/src/main/java/org/apache/gravitino/cache/BaseEntityCache.java
index e23c53b81d..23143bf331 100644
--- a/core/src/main/java/org/apache/gravitino/cache/BaseEntityCache.java
+++ b/core/src/main/java/org/apache/gravitino/cache/BaseEntityCache.java
@@ -49,8 +49,8 @@ public abstract class BaseEntityCache implements EntityCache {
* (an old comment, property, or job status), never a wrong pointer, and
each can be invalidated
* with a single one-to-one key drop. Every other type is read straight from
the store.
* User/group/role embed relation-derived data that a per-node cache cannot
invalidate;
- * model/model version and function carry a load-bearing pointer that would
be silently wrong if
- * served stale.
+ * model/model version, Semantic Model, and function carry load-bearing
content that would be
+ * silently wrong if served stale.
*/
private static final Set<Entity.EntityType> CACHEABLE_TYPES =
Sets.immutableEnumSet(
diff --git
a/core/src/main/java/org/apache/gravitino/cache/SupportsEntityStoreCache.java
b/core/src/main/java/org/apache/gravitino/cache/SupportsEntityStoreCache.java
index a6eb72d11a..2f46846091 100644
---
a/core/src/main/java/org/apache/gravitino/cache/SupportsEntityStoreCache.java
+++
b/core/src/main/java/org/apache/gravitino/cache/SupportsEntityStoreCache.java
@@ -70,9 +70,9 @@ public interface SupportsEntityStoreCache {
*
* <p>Only self-contained entities are approved. User/group/role are
deliberately excluded because
* they contain relation-derived data whose source can change through
another entity. Model/model
- * version and function are also excluded because they carry load-bearing
pointers that must not
- * be served stale. Other unapproved and newly introduced entity types go
straight to the store
- * until their invalidation behavior has been validated.
+ * version, Semantic Model, and function are also excluded because they
carry load-bearing content
+ * that must not be served stale. Other unapproved and newly introduced
entity types go straight
+ * to the store until their invalidation behavior has been validated.
*
* <p>Implementations must also invoke {@link
#invalidateOnKeyChange(Entity)} for every entity,
* including non-cacheable ones, since a non-cacheable entity may still
invalidate a cacheable
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
b/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
index 7ab84f9005..c847c8f910 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
@@ -63,6 +63,7 @@ import org.apache.gravitino.meta.ModelVersionEntity;
import org.apache.gravitino.meta.PolicyEntity;
import org.apache.gravitino.meta.RoleEntity;
import org.apache.gravitino.meta.SchemaEntity;
+import org.apache.gravitino.meta.SemanticModelEntity;
import org.apache.gravitino.meta.StatisticEntity;
import org.apache.gravitino.meta.TableEntity;
import org.apache.gravitino.meta.TagEntity;
@@ -88,6 +89,7 @@ import
org.apache.gravitino.storage.relational.service.PolicyMetaService;
import org.apache.gravitino.storage.relational.service.PolicyTagRelService;
import org.apache.gravitino.storage.relational.service.RoleMetaService;
import org.apache.gravitino.storage.relational.service.SchemaMetaService;
+import
org.apache.gravitino.storage.relational.service.SemanticModelMetaService;
import org.apache.gravitino.storage.relational.service.StatisticMetaService;
import org.apache.gravitino.storage.relational.service.TableColumnMetaService;
import org.apache.gravitino.storage.relational.service.TableMetaService;
@@ -281,6 +283,8 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
return (E) JobMetaService.getInstance().getJobByIdentifier(ident);
case VIEW:
return (E) ViewMetaService.getInstance().getViewByIdentifier(ident);
+ case SEMANTIC_MODEL:
+ return (E)
SemanticModelMetaService.getInstance().getSemanticModelByIdentifier(ident);
default:
throw new UnsupportedEntityTypeException(
"Unsupported entity type: %s for get operation", entityType);
@@ -993,6 +997,9 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
JobMetaService.getInstance().insertJob((JobEntity) e, overwritten);
} else if (e instanceof ViewEntity) {
ViewMetaService.getInstance().insertView((ViewEntity) e, overwritten);
+ } else if (e instanceof SemanticModelEntity) {
+ SemanticModelMetaService.getInstance()
+ .insertSemanticModel((SemanticModelEntity) e, overwritten);
} else if (e instanceof GenericEntity) {
GenericEntity genericEntity = (GenericEntity) e;
throw new UnsupportedEntityTypeException(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStoreIdResolver.java
b/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStoreIdResolver.java
index 895572719a..e5a4aae6fb 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStoreIdResolver.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStoreIdResolver.java
@@ -37,6 +37,7 @@ import
org.apache.gravitino.storage.relational.service.ModelMetaService;
import org.apache.gravitino.storage.relational.service.PolicyMetaService;
import org.apache.gravitino.storage.relational.service.RoleMetaService;
import org.apache.gravitino.storage.relational.service.SchemaMetaService;
+import
org.apache.gravitino.storage.relational.service.SemanticModelMetaService;
import org.apache.gravitino.storage.relational.service.TableColumnMetaService;
import org.apache.gravitino.storage.relational.service.TableMetaService;
import org.apache.gravitino.storage.relational.service.TagMetaService;
@@ -69,7 +70,8 @@ public class RelationalEntityStoreIdResolver implements
EntityIdResolver {
Entity.EntityType.MODEL,
Entity.EntityType.COLUMN,
Entity.EntityType.FUNCTION,
- Entity.EntityType.VIEW);
+ Entity.EntityType.VIEW,
+ Entity.EntityType.SEMANTIC_MODEL);
@Override
public NamespacedEntityId getEntityIds(NameIdentifier nameIdentifier,
Entity.EntityType type) {
@@ -235,6 +237,17 @@ public class RelationalEntityStoreIdResolver implements
EntityIdResolver {
return new NamespacedEntityId(
viewId, schemaIds.getMetalakeId(), schemaIds.getCatalogId(),
schemaIds.getSchemaId());
+ case SEMANTIC_MODEL:
+ long semanticModelId =
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelIdBySchemaIdAndName(
+ schemaIds.getSchemaId(), nameIdentifier.name());
+ return new NamespacedEntityId(
+ semanticModelId,
+ schemaIds.getMetalakeId(),
+ schemaIds.getCatalogId(),
+ schemaIds.getSchemaId());
+
case FUNCTION:
long functionId =
FunctionMetaService.getInstance()
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelMetaMapper.java
new file mode 100644
index 0000000000..a8690f8230
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelMetaMapper.java
@@ -0,0 +1,143 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper;
+
+import org.apache.gravitino.storage.relational.po.SemanticModelPO;
+import org.apache.gravitino.storage.relational.po.SemanticModelVersionInfoPO;
+import org.apache.ibatis.annotations.InsertProvider;
+import org.apache.ibatis.annotations.One;
+import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.annotations.Result;
+import org.apache.ibatis.annotations.ResultMap;
+import org.apache.ibatis.annotations.Results;
+import org.apache.ibatis.annotations.Select;
+import org.apache.ibatis.annotations.SelectProvider;
+import org.apache.ibatis.annotations.UpdateProvider;
+
+/** A MyBatis mapper for Semantic Model create and load operations. */
+public interface SemanticModelMetaMapper {
+
+ /** The Semantic Model identity table name. */
+ String TABLE_NAME = "semantic_model_meta";
+
+ /** The Semantic Model version snapshot table name. */
+ String VERSION_TABLE_NAME = "semantic_model_version_info";
+
+ /** Declares the nested result mapping for a Semantic Model version
snapshot. */
+ @Results(
+ id = "mapToSemanticModelVersionInfoPO",
+ value = {
+ @Result(property = "id", column = "id", id = true),
+ @Result(property = "metalakeId", column = "version_metalake_id"),
+ @Result(property = "catalogId", column = "version_catalog_id"),
+ @Result(property = "schemaId", column = "version_schema_id"),
+ @Result(property = "semanticModelId", column =
"version_semantic_model_id"),
+ @Result(property = "version", column = "version"),
+ @Result(property = "semanticModelName", column =
"version_semantic_model_name"),
+ @Result(property = "semanticModelComment", column =
"semantic_model_comment"),
+ @Result(property = "semanticModelDefinition", column =
"semantic_model_definition"),
+ @Result(property = "properties", column = "properties"),
+ @Result(property = "auditInfo", column = "version_audit_info"),
+ @Result(property = "deletedAt", column = "version_deleted_at")
+ })
+ @Select("SELECT 1")
+ SemanticModelVersionInfoPO mapToSemanticModelVersionInfoPO();
+
+ /** Selects an active Semantic Model ID by schema ID and name. */
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "selectSemanticModelIdBySchemaIdAndName")
+ Long selectSemanticModelIdBySchemaIdAndName(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName);
+
+ /** Selects a current Semantic Model snapshot by schema ID and name. */
+ @Results(
+ id = "semanticModelPOResultMap",
+ value = {
+ @Result(property = "semanticModelId", column = "semantic_model_id", id
= true),
+ @Result(property = "semanticModelName", column =
"semantic_model_name"),
+ @Result(property = "metalakeId", column = "metalake_id"),
+ @Result(property = "catalogId", column = "catalog_id"),
+ @Result(property = "schemaId", column = "schema_id"),
+ @Result(property = "currentVersion", column = "current_version"),
+ @Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "auditInfo", column = "audit_info"),
+ @Result(property = "deletedAt", column = "deleted_at"),
+ @Result(
+ property = "semanticModelVersionInfoPO",
+ javaType = SemanticModelVersionInfoPO.class,
+ column =
+ "{id,version_metalake_id,version_catalog_id,version_schema_id,"
+ +
"version_semantic_model_id,version,version_semantic_model_name,"
+ +
"semantic_model_comment,semantic_model_definition,properties,"
+ + "version_audit_info,version_deleted_at}",
+ one = @One(resultMap = "mapToSemanticModelVersionInfoPO"))
+ })
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "selectSemanticModelMetaBySchemaIdAndName")
+ SemanticModelPO selectSemanticModelMetaBySchemaIdAndName(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName);
+
+ /** Selects and locks an active Semantic Model identity by schema ID and
name. */
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "selectSemanticModelMetaBySchemaIdAndNameForUpdate")
+ SemanticModelPO selectSemanticModelMetaBySchemaIdAndNameForUpdate(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName);
+
+ /** Selects a current Semantic Model snapshot by stable ID. */
+ @ResultMap("semanticModelPOResultMap")
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "selectSemanticModelMetaById")
+ SemanticModelPO selectSemanticModelMetaById(@Param("semanticModelId") Long
semanticModelId);
+
+ /** Selects and locks a Semantic Model identity by stable ID. */
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "selectSemanticModelMetaByIdForUpdate")
+ SemanticModelPO selectSemanticModelMetaByIdForUpdate(
+ @Param("semanticModelId") Long semanticModelId);
+
+ /** Selects a current Semantic Model snapshot by fully qualified name. */
+ @ResultMap("semanticModelPOResultMap")
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "selectSemanticModelByFullQualifiedName")
+ SemanticModelPO selectSemanticModelByFullQualifiedName(
+ @Param("metalakeName") String metalakeName,
+ @Param("catalogName") String catalogName,
+ @Param("schemaName") String schemaName,
+ @Param("semanticModelName") String semanticModelName);
+
+ /** Inserts a Semantic Model identity row. */
+ @InsertProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "insertSemanticModelMeta")
+ void insertSemanticModelMeta(@Param("semanticModelMeta") SemanticModelPO
semanticModelPO);
+
+ /** Updates a Semantic Model identity while its current version is
unchanged. */
+ @UpdateProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "updateSemanticModelMeta")
+ Integer updateSemanticModelMeta(
+ @Param("newSemanticModelMeta") SemanticModelPO newSemanticModelPO,
+ @Param("oldSemanticModelMeta") SemanticModelPO oldSemanticModelPO);
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelMetaSQLProviderFactory.java
new file mode 100644
index 0000000000..ca3e3d5865
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelMetaSQLProviderFactory.java
@@ -0,0 +1,107 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper;
+
+import com.google.common.collect.ImmutableMap;
+import java.util.Map;
+import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType;
+import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelMetaBaseSQLProvider;
+import
org.apache.gravitino.storage.relational.mapper.provider.postgresql.SemanticModelMetaPostgreSQLProvider;
+import org.apache.gravitino.storage.relational.po.SemanticModelPO;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.ibatis.annotations.Param;
+
+/** Selects database-specific SQL providers for Semantic Model create and load
operations. */
+public class SemanticModelMetaSQLProviderFactory {
+
+ private static final Map<JDBCBackendType, SemanticModelMetaBaseSQLProvider>
+ SEMANTIC_MODEL_META_SQL_PROVIDER_MAP =
+ ImmutableMap.of(
+ JDBCBackendType.MYSQL, new SemanticModelMetaMySQLProvider(),
+ JDBCBackendType.H2, new SemanticModelMetaH2Provider(),
+ JDBCBackendType.POSTGRESQL, new
SemanticModelMetaPostgreSQLProvider());
+
+ /** Returns the SQL provider for the configured relational backend. */
+ public static SemanticModelMetaBaseSQLProvider getProvider() {
+ String databaseId =
+ SqlSessionFactoryHelper.getInstance()
+ .getSqlSessionFactory()
+ .getConfiguration()
+ .getDatabaseId();
+ return
SEMANTIC_MODEL_META_SQL_PROVIDER_MAP.get(JDBCBackendType.fromString(databaseId));
+ }
+
+ static class SemanticModelMetaMySQLProvider extends
SemanticModelMetaBaseSQLProvider {}
+
+ static class SemanticModelMetaH2Provider extends
SemanticModelMetaBaseSQLProvider {}
+
+ /** Provides SQL for selecting a Semantic Model ID by schema ID and name. */
+ public static String selectSemanticModelIdBySchemaIdAndName(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName) {
+ return getProvider().selectSemanticModelIdBySchemaIdAndName(schemaId,
semanticModelName);
+ }
+
+ /** Provides SQL for selecting a Semantic Model by schema ID and name. */
+ public static String selectSemanticModelMetaBySchemaIdAndName(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName) {
+ return getProvider().selectSemanticModelMetaBySchemaIdAndName(schemaId,
semanticModelName);
+ }
+
+ /** Provides SQL for selecting and locking a Semantic Model identity by
schema ID and name. */
+ public static String selectSemanticModelMetaBySchemaIdAndNameForUpdate(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName) {
+ return getProvider()
+ .selectSemanticModelMetaBySchemaIdAndNameForUpdate(schemaId,
semanticModelName);
+ }
+
+ /** Provides SQL for selecting a Semantic Model by stable ID. */
+ public static String selectSemanticModelMetaById(@Param("semanticModelId")
Long semanticModelId) {
+ return getProvider().selectSemanticModelMetaById(semanticModelId);
+ }
+
+ /** Provides SQL for selecting and locking a Semantic Model identity by
stable ID. */
+ public static String selectSemanticModelMetaByIdForUpdate(
+ @Param("semanticModelId") Long semanticModelId) {
+ return getProvider().selectSemanticModelMetaByIdForUpdate(semanticModelId);
+ }
+
+ /** Provides SQL for selecting a Semantic Model by fully qualified name. */
+ public static String selectSemanticModelByFullQualifiedName(
+ @Param("metalakeName") String metalakeName,
+ @Param("catalogName") String catalogName,
+ @Param("schemaName") String schemaName,
+ @Param("semanticModelName") String semanticModelName) {
+ return getProvider()
+ .selectSemanticModelByFullQualifiedName(
+ metalakeName, catalogName, schemaName, semanticModelName);
+ }
+
+ /** Provides SQL for inserting a Semantic Model identity. */
+ public static String insertSemanticModelMeta(
+ @Param("semanticModelMeta") SemanticModelPO semanticModelPO) {
+ return getProvider().insertSemanticModelMeta(semanticModelPO);
+ }
+
+ /** Provides SQL for updating a Semantic Model identity when its version is
unchanged. */
+ public static String updateSemanticModelMeta(
+ @Param("newSemanticModelMeta") SemanticModelPO newSemanticModelPO,
+ @Param("oldSemanticModelMeta") SemanticModelPO oldSemanticModelPO) {
+ return getProvider().updateSemanticModelMeta(newSemanticModelPO,
oldSemanticModelPO);
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelVersionInfoMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelVersionInfoMapper.java
new file mode 100644
index 0000000000..c02b2c43be
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelVersionInfoMapper.java
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper;
+
+import org.apache.gravitino.storage.relational.po.SemanticModelVersionInfoPO;
+import org.apache.ibatis.annotations.InsertProvider;
+import org.apache.ibatis.annotations.Param;
+
+/** A MyBatis mapper for creating Semantic Model version snapshots. */
+public interface SemanticModelVersionInfoMapper {
+
+ /** The Semantic Model version snapshot table name. */
+ String TABLE_NAME = "semantic_model_version_info";
+
+ /** Inserts a Semantic Model version snapshot. */
+ @InsertProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "insertSemanticModelVersionInfo")
+ void insertSemanticModelVersionInfo(
+ @Param("semanticModelVersionInfo") SemanticModelVersionInfoPO
versionInfoPO);
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelVersionInfoSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelVersionInfoSQLProviderFactory.java
new file mode 100644
index 0000000000..7e61cfc110
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SemanticModelVersionInfoSQLProviderFactory.java
@@ -0,0 +1,60 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper;
+
+import com.google.common.collect.ImmutableMap;
+import java.util.Map;
+import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType;
+import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelVersionInfoBaseSQLProvider;
+import
org.apache.gravitino.storage.relational.mapper.provider.postgresql.SemanticModelVersionInfoPostgreSQLProvider;
+import org.apache.gravitino.storage.relational.po.SemanticModelVersionInfoPO;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.ibatis.annotations.Param;
+
+/** Selects database-specific SQL providers for Semantic Model snapshot
creation. */
+public class SemanticModelVersionInfoSQLProviderFactory {
+
+ private static final Map<JDBCBackendType,
SemanticModelVersionInfoBaseSQLProvider>
+ SEMANTIC_MODEL_VERSION_INFO_SQL_PROVIDER_MAP =
+ ImmutableMap.of(
+ JDBCBackendType.MYSQL, new
SemanticModelVersionInfoMySQLProvider(),
+ JDBCBackendType.H2, new SemanticModelVersionInfoH2Provider(),
+ JDBCBackendType.POSTGRESQL, new
SemanticModelVersionInfoPostgreSQLProvider());
+
+ /** Returns the SQL provider for the configured relational backend. */
+ public static SemanticModelVersionInfoBaseSQLProvider getProvider() {
+ String databaseId =
+ SqlSessionFactoryHelper.getInstance()
+ .getSqlSessionFactory()
+ .getConfiguration()
+ .getDatabaseId();
+ return
SEMANTIC_MODEL_VERSION_INFO_SQL_PROVIDER_MAP.get(JDBCBackendType.fromString(databaseId));
+ }
+
+ static class SemanticModelVersionInfoMySQLProvider
+ extends SemanticModelVersionInfoBaseSQLProvider {}
+
+ static class SemanticModelVersionInfoH2Provider extends
SemanticModelVersionInfoBaseSQLProvider {}
+
+ /** Provides SQL for inserting a Semantic Model version snapshot. */
+ public static String insertSemanticModelVersionInfo(
+ @Param("semanticModelVersionInfo") SemanticModelVersionInfoPO
versionInfoPO) {
+ return getProvider().insertSemanticModelVersionInfo(versionInfoPO);
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/DefaultMapperPackageProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/DefaultMapperPackageProvider.java
index 6a2bd9ca85..c05336d88c 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/DefaultMapperPackageProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/DefaultMapperPackageProvider.java
@@ -42,6 +42,8 @@ import
org.apache.gravitino.storage.relational.mapper.PolicyVersionMapper;
import org.apache.gravitino.storage.relational.mapper.RoleMetaMapper;
import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
import org.apache.gravitino.storage.relational.mapper.SecurableObjectMapper;
+import org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper;
+import
org.apache.gravitino.storage.relational.mapper.SemanticModelVersionInfoMapper;
import org.apache.gravitino.storage.relational.mapper.StatisticMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TableColumnMapper;
import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
@@ -82,6 +84,8 @@ public class DefaultMapperPackageProvider implements
MapperPackageProvider {
RoleMetaMapper.class,
SchemaMetaMapper.class,
SecurableObjectMapper.class,
+ SemanticModelMetaMapper.class,
+ SemanticModelVersionInfoMapper.class,
StatisticMetaMapper.class,
TableColumnMapper.class,
TableMetaMapper.class,
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SemanticModelMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SemanticModelMetaBaseSQLProvider.java
new file mode 100644
index 0000000000..f32a439959
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SemanticModelMetaBaseSQLProvider.java
@@ -0,0 +1,208 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper.provider.base;
+
+import static
org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper.TABLE_NAME;
+import static
org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper.VERSION_TABLE_NAME;
+
+import org.apache.gravitino.storage.relational.mapper.CatalogMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.MetalakeMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
+import org.apache.gravitino.storage.relational.po.SemanticModelPO;
+import org.apache.ibatis.annotations.Param;
+
+/** Provides MySQL-compatible SQL for Semantic Model create and load
operations. */
+public class SemanticModelMetaBaseSQLProvider {
+
+ private static final String CURRENT_SNAPSHOT_COLUMNS =
+ " smm.semantic_model_id, smm.semantic_model_name, smm.metalake_id,"
+ + " smm.catalog_id, smm.schema_id, smm.current_version,
smm.last_version,"
+ + " smm.audit_info, smm.deleted_at, smvi.id,"
+ + " smvi.metalake_id as version_metalake_id,"
+ + " smvi.catalog_id as version_catalog_id,"
+ + " smvi.schema_id as version_schema_id,"
+ + " smvi.semantic_model_id as version_semantic_model_id,
smvi.version,"
+ + " smvi.semantic_model_name as version_semantic_model_name,"
+ + " smvi.semantic_model_comment, smvi.semantic_model_definition,
smvi.properties,"
+ + " smvi.audit_info as version_audit_info,"
+ + " smvi.deleted_at as version_deleted_at";
+
+ /** Returns SQL for selecting an active Semantic Model ID by schema ID and
name. */
+ public String selectSemanticModelIdBySchemaIdAndName(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName) {
+ return "SELECT semantic_model_id as semanticModelId FROM "
+ + TABLE_NAME
+ + " WHERE schema_id = #{schemaId}"
+ + " AND semantic_model_name = #{semanticModelName} AND deleted_at = 0";
+ }
+
+ /** Returns SQL for selecting a current Semantic Model snapshot by schema ID
and name. */
+ public String selectSemanticModelMetaBySchemaIdAndName(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName) {
+ return "SELECT"
+ + CURRENT_SNAPSHOT_COLUMNS
+ + " FROM "
+ + TABLE_NAME
+ + " smm INNER JOIN "
+ + VERSION_TABLE_NAME
+ + " smvi ON smm.semantic_model_id = smvi.semantic_model_id"
+ + " AND smm.current_version = smvi.version"
+ + " WHERE smm.schema_id = #{schemaId}"
+ + " AND smm.semantic_model_name = #{semanticModelName}"
+ + " AND smm.deleted_at = 0 AND smvi.deleted_at = 0";
+ }
+
+ /** Returns SQL for selecting and locking an active Semantic Model identity
by natural key. */
+ public String selectSemanticModelMetaBySchemaIdAndNameForUpdate(
+ @Param("schemaId") Long schemaId, @Param("semanticModelName") String
semanticModelName) {
+ return "SELECT semantic_model_id as semanticModelId,"
+ + " semantic_model_name as semanticModelName, metalake_id as
metalakeId,"
+ + " catalog_id as catalogId, schema_id as schemaId,"
+ + " current_version as currentVersion, last_version as lastVersion,"
+ + " audit_info as auditInfo, deleted_at as deletedAt"
+ + " FROM "
+ + TABLE_NAME
+ + " WHERE schema_id = #{schemaId}"
+ + " AND semantic_model_name = #{semanticModelName}"
+ + " AND deleted_at = 0 FOR UPDATE";
+ }
+
+ /** Returns SQL for selecting a current Semantic Model snapshot by stable
ID. */
+ public String selectSemanticModelMetaById(@Param("semanticModelId") Long
semanticModelId) {
+ return "SELECT"
+ + CURRENT_SNAPSHOT_COLUMNS
+ + " FROM "
+ + TABLE_NAME
+ + " smm INNER JOIN "
+ + VERSION_TABLE_NAME
+ + " smvi ON smm.semantic_model_id = smvi.semantic_model_id"
+ + " AND smm.current_version = smvi.version"
+ + " WHERE smm.semantic_model_id = #{semanticModelId}"
+ + " AND smm.deleted_at = 0 AND smvi.deleted_at = 0";
+ }
+
+ /** Returns SQL for selecting and locking a Semantic Model identity by
stable ID. */
+ public String selectSemanticModelMetaByIdForUpdate(
+ @Param("semanticModelId") Long semanticModelId) {
+ return "SELECT semantic_model_id as semanticModelId,"
+ + " semantic_model_name as semanticModelName, metalake_id as
metalakeId,"
+ + " catalog_id as catalogId, schema_id as schemaId,"
+ + " current_version as currentVersion, last_version as lastVersion,"
+ + " audit_info as auditInfo, deleted_at as deletedAt FROM "
+ + TABLE_NAME
+ + " WHERE semantic_model_id = #{semanticModelId} AND deleted_at = 0
FOR UPDATE";
+ }
+
+ /** Returns SQL for selecting a current Semantic Model by fully qualified
name. */
+ public String selectSemanticModelByFullQualifiedName(
+ @Param("metalakeName") String metalakeName,
+ @Param("catalogName") String catalogName,
+ @Param("schemaName") String schemaName,
+ @Param("semanticModelName") String semanticModelName) {
+ return """
+ SELECT
+ mm.metalake_id,
+ cm.catalog_id,
+ sm.schema_id,
+ smm.semantic_model_id,
+ smm.semantic_model_name,
+ smm.current_version,
+ smm.last_version,
+ smm.audit_info,
+ smm.deleted_at,
+ smvi.id,
+ smvi.metalake_id as version_metalake_id,
+ smvi.catalog_id as version_catalog_id,
+ smvi.schema_id as version_schema_id,
+ smvi.semantic_model_id as version_semantic_model_id,
+ smvi.version,
+ smvi.semantic_model_name as version_semantic_model_name,
+ smvi.semantic_model_comment,
+ smvi.semantic_model_definition,
+ smvi.properties,
+ smvi.audit_info as version_audit_info,
+ smvi.deleted_at as version_deleted_at
+ FROM
+ %s mm
+ INNER JOIN
+ %s cm ON mm.metalake_id = cm.metalake_id
+ AND cm.catalog_name = #{catalogName}
+ AND cm.deleted_at = 0
+ LEFT JOIN
+ %s sm ON cm.catalog_id = sm.catalog_id
+ AND sm.schema_name = #{schemaName}
+ AND sm.deleted_at = 0
+ LEFT JOIN
+ %s smm ON sm.schema_id = smm.schema_id
+ AND smm.semantic_model_name = #{semanticModelName}
+ AND smm.deleted_at = 0
+ LEFT JOIN
+ %s smvi ON smm.semantic_model_id = smvi.semantic_model_id
+ AND smm.current_version = smvi.version
+ AND smvi.deleted_at = 0
+ WHERE
+ mm.metalake_name = #{metalakeName}
+ AND mm.deleted_at = 0
+ """
+ .formatted(
+ MetalakeMetaMapper.TABLE_NAME,
+ CatalogMetaMapper.TABLE_NAME,
+ SchemaMetaMapper.TABLE_NAME,
+ TABLE_NAME,
+ VERSION_TABLE_NAME);
+ }
+
+ /** Returns SQL for inserting a Semantic Model identity row. */
+ public String insertSemanticModelMeta(
+ @Param("semanticModelMeta") SemanticModelPO semanticModelPO) {
+ return "INSERT INTO "
+ + TABLE_NAME
+ + " (semantic_model_id, semantic_model_name, metalake_id, catalog_id,
schema_id,"
+ + " current_version, last_version, audit_info, deleted_at)"
+ + " VALUES (#{semanticModelMeta.semanticModelId},"
+ + " #{semanticModelMeta.semanticModelName},
#{semanticModelMeta.metalakeId},"
+ + " #{semanticModelMeta.catalogId}, #{semanticModelMeta.schemaId},"
+ + " #{semanticModelMeta.currentVersion},
#{semanticModelMeta.lastVersion},"
+ + " #{semanticModelMeta.auditInfo}, #{semanticModelMeta.deletedAt})";
+ }
+
+ /** Returns SQL for updating a Semantic Model identity with a version check.
*/
+ public String updateSemanticModelMeta(
+ @Param("newSemanticModelMeta") SemanticModelPO newSemanticModelPO,
+ @Param("oldSemanticModelMeta") SemanticModelPO oldSemanticModelPO) {
+ return "UPDATE "
+ + TABLE_NAME
+ + " SET semantic_model_name =
#{newSemanticModelMeta.semanticModelName},"
+ + " metalake_id = #{newSemanticModelMeta.metalakeId},"
+ + " catalog_id = #{newSemanticModelMeta.catalogId},"
+ + " schema_id = #{newSemanticModelMeta.schemaId},"
+ + " current_version = #{newSemanticModelMeta.currentVersion},"
+ + " last_version = #{newSemanticModelMeta.lastVersion},"
+ + " audit_info = #{newSemanticModelMeta.auditInfo},"
+ + " deleted_at = #{newSemanticModelMeta.deletedAt}"
+ + " WHERE semantic_model_id = #{oldSemanticModelMeta.semanticModelId}"
+ + " AND current_version = #{oldSemanticModelMeta.currentVersion}"
+ + " AND deleted_at = 0"
+ + " AND NOT EXISTS (SELECT 1 FROM "
+ + VERSION_TABLE_NAME
+ + " smvi WHERE smvi.semantic_model_id =
#{oldSemanticModelMeta.semanticModelId}"
+ + " AND smvi.version >= #{newSemanticModelMeta.currentVersion}"
+ + " AND smvi.deleted_at = 0)";
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SemanticModelVersionInfoBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SemanticModelVersionInfoBaseSQLProvider.java
new file mode 100644
index 0000000000..8873f1c77e
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SemanticModelVersionInfoBaseSQLProvider.java
@@ -0,0 +1,45 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper.provider.base;
+
+import
org.apache.gravitino.storage.relational.mapper.SemanticModelVersionInfoMapper;
+import org.apache.gravitino.storage.relational.po.SemanticModelVersionInfoPO;
+import org.apache.ibatis.annotations.Param;
+
+/** Provides MySQL-compatible SQL for creating Semantic Model version
snapshots. */
+public class SemanticModelVersionInfoBaseSQLProvider {
+
+ /** Returns SQL for inserting a Semantic Model version snapshot. */
+ public String insertSemanticModelVersionInfo(
+ @Param("semanticModelVersionInfo") SemanticModelVersionInfoPO
versionInfoPO) {
+ return "INSERT INTO "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " (metalake_id, catalog_id, schema_id, semantic_model_id, version,"
+ + " semantic_model_name, semantic_model_comment,
semantic_model_definition,"
+ + " properties, audit_info, deleted_at)"
+ + " VALUES (#{semanticModelVersionInfo.metalakeId},"
+ + " #{semanticModelVersionInfo.catalogId},
#{semanticModelVersionInfo.schemaId},"
+ + " #{semanticModelVersionInfo.semanticModelId},
#{semanticModelVersionInfo.version},"
+ + " #{semanticModelVersionInfo.semanticModelName},"
+ + " #{semanticModelVersionInfo.semanticModelComment},"
+ + " #{semanticModelVersionInfo.semanticModelDefinition},"
+ + " #{semanticModelVersionInfo.properties},
#{semanticModelVersionInfo.auditInfo},"
+ + " #{semanticModelVersionInfo.deletedAt})";
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SemanticModelMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SemanticModelMetaPostgreSQLProvider.java
new file mode 100644
index 0000000000..1261cfbcce
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SemanticModelMetaPostgreSQLProvider.java
@@ -0,0 +1,24 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper.provider.postgresql;
+
+import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelMetaBaseSQLProvider;
+
+/** Provides PostgreSQL SQL for Semantic Model create and load operations. */
+public class SemanticModelMetaPostgreSQLProvider extends
SemanticModelMetaBaseSQLProvider {}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SemanticModelVersionInfoPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SemanticModelVersionInfoPostgreSQLProvider.java
new file mode 100644
index 0000000000..9f78550b93
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SemanticModelVersionInfoPostgreSQLProvider.java
@@ -0,0 +1,25 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.mapper.provider.postgresql;
+
+import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelVersionInfoBaseSQLProvider;
+
+/** Provides PostgreSQL SQL for creating Semantic Model version snapshots. */
+public class SemanticModelVersionInfoPostgreSQLProvider
+ extends SemanticModelVersionInfoBaseSQLProvider {}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/po/SemanticModelPO.java
b/core/src/main/java/org/apache/gravitino/storage/relational/po/SemanticModelPO.java
new file mode 100644
index 0000000000..34e95faaab
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/po/SemanticModelPO.java
@@ -0,0 +1,228 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.po;
+
+import static
org.apache.gravitino.storage.relational.utils.POConverters.DEFAULT_DELETED_AT;
+
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.google.common.base.Preconditions;
+import java.util.Collections;
+import java.util.Map;
+import lombok.EqualsAndHashCode;
+import lombok.Getter;
+import lombok.ToString;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.Namespace;
+import org.apache.gravitino.dto.semantic.SemanticModelDefinitionDTO;
+import org.apache.gravitino.json.JsonUtils;
+import org.apache.gravitino.meta.AuditInfo;
+import org.apache.gravitino.meta.NamespacedEntityId;
+import org.apache.gravitino.meta.SemanticModelEntity;
+import org.apache.gravitino.storage.relational.service.EntityIdService;
+
+/** The persistent object for Semantic Model identity metadata and its current
version snapshot. */
+@Getter
+@EqualsAndHashCode(exclude = "semanticModelVersionInfoPO")
+@ToString
+public class SemanticModelPO {
+
+ /** The initial version allocated to a newly created Semantic Model. */
+ public static final Integer INITIAL_VERSION = 1;
+
+ private Long semanticModelId;
+ private String semanticModelName;
+ private Long metalakeId;
+ private Long catalogId;
+ private Long schemaId;
+ private String auditInfo;
+ private Integer currentVersion;
+ private Integer lastVersion;
+ private Long deletedAt;
+ private SemanticModelVersionInfoPO semanticModelVersionInfoPO;
+
+ /** Creates an empty persistent object for MyBatis. */
+ public SemanticModelPO() {}
+
+ /** A Lombok builder for {@link SemanticModelPO}. */
+ public static class SemanticModelPOBuilder {
+ // Lombok generates the builder methods.
+ }
+
+ @lombok.Builder(setterPrefix = "with")
+ private SemanticModelPO(
+ Long semanticModelId,
+ String semanticModelName,
+ Long metalakeId,
+ Long catalogId,
+ Long schemaId,
+ String auditInfo,
+ Integer currentVersion,
+ Integer lastVersion,
+ Long deletedAt,
+ SemanticModelVersionInfoPO semanticModelVersionInfoPO) {
+ Preconditions.checkArgument(semanticModelId != null, "Semantic Model id is
required");
+ Preconditions.checkArgument(semanticModelName != null, "Semantic Model
name is required");
+ Preconditions.checkArgument(metalakeId != null, "Metalake id is required");
+ Preconditions.checkArgument(catalogId != null, "Catalog id is required");
+ Preconditions.checkArgument(schemaId != null, "Schema id is required");
+ Preconditions.checkArgument(auditInfo != null, "Audit info is required");
+ Preconditions.checkArgument(currentVersion != null, "Current version is
required");
+ Preconditions.checkArgument(lastVersion != null, "Last version is
required");
+ Preconditions.checkArgument(deletedAt != null, "Deleted at is required");
+
+ this.semanticModelId = semanticModelId;
+ this.semanticModelName = semanticModelName;
+ this.metalakeId = metalakeId;
+ this.catalogId = catalogId;
+ this.schemaId = schemaId;
+ this.auditInfo = auditInfo;
+ this.currentVersion = currentVersion;
+ this.lastVersion = lastVersion;
+ this.deletedAt = deletedAt;
+ this.semanticModelVersionInfoPO = semanticModelVersionInfoPO;
+ }
+
+ /**
+ * Converts a persistent object and its current version snapshot to a
Semantic Model entity.
+ *
+ * @param semanticModelPO The persistent object to convert.
+ * @param namespace The Semantic Model namespace.
+ * @return The converted Semantic Model entity.
+ */
+ public static SemanticModelEntity fromSemanticModelPO(
+ SemanticModelPO semanticModelPO, Namespace namespace) {
+ try {
+ SemanticModelVersionInfoPO versionPO =
semanticModelPO.getSemanticModelVersionInfoPO();
+ SemanticModelDefinitionDTO definitionDTO =
+ JsonUtils.anyFieldMapper()
+ .readValue(versionPO.semanticModelDefinition(),
SemanticModelDefinitionDTO.class);
+ Map<String, String> properties =
+ versionPO.properties() == null
+ ? Collections.emptyMap()
+ : JsonUtils.anyFieldMapper()
+ .readValue(
+ versionPO.properties(),
+ JsonUtils.anyFieldMapper()
+ .getTypeFactory()
+ .constructMapType(Map.class, String.class,
String.class));
+
+ return SemanticModelEntity.builder()
+ .withId(semanticModelPO.getSemanticModelId())
+ .withName(versionPO.semanticModelName())
+ .withNamespace(namespace)
+ .withComment(versionPO.semanticModelComment())
+ .withDefinition(definitionDTO.toDefinition())
+ .withProperties(properties)
+ .withAuditInfo(
+
JsonUtils.anyFieldMapper().readValue(semanticModelPO.getAuditInfo(),
AuditInfo.class))
+ .build();
+ } catch (JsonProcessingException e) {
+ throw new RuntimeException("Failed to deserialize Semantic Model JSON",
e);
+ }
+ }
+
+ /**
+ * Initializes a new Semantic Model identity and version-one snapshot.
+ *
+ * @param semanticModelEntity The Semantic Model entity.
+ * @param builder The identity persistent-object builder.
+ * @return The initialized persistent object.
+ */
+ public static SemanticModelPO initializeSemanticModelPO(
+ SemanticModelEntity semanticModelEntity, SemanticModelPOBuilder builder)
{
+ return buildSemanticModelPO(semanticModelEntity, builder, INITIAL_VERSION);
+ }
+
+ /**
+ * Creates a complete version snapshot for a Semantic Model entity.
+ *
+ * @param semanticModelEntity The Semantic Model entity.
+ * @param namespacedEntityId The resolved schema and ancestor IDs.
+ * @param version The version to allocate.
+ * @return The version snapshot persistent object.
+ */
+ public static SemanticModelVersionInfoPO
initializeSemanticModelVersionInfoPO(
+ SemanticModelEntity semanticModelEntity,
+ NamespacedEntityId namespacedEntityId,
+ Integer version) {
+ try {
+ String definitionJson =
+ JsonUtils.anyFieldMapper()
+ .writeValueAsString(
+
SemanticModelDefinitionDTO.fromDefinition(semanticModelEntity.definition()));
+ String propertiesJson =
+ semanticModelEntity.properties().isEmpty()
+ ? null
+ :
JsonUtils.anyFieldMapper().writeValueAsString(semanticModelEntity.properties());
+
+ return SemanticModelVersionInfoPO.builder()
+ .withSemanticModelId(semanticModelEntity.id())
+ .withMetalakeId(namespacedEntityId.namespaceIds()[0])
+ .withCatalogId(namespacedEntityId.namespaceIds()[1])
+ .withSchemaId(namespacedEntityId.entityId())
+ .withVersion(version)
+ .withSemanticModelName(semanticModelEntity.name())
+ .withSemanticModelComment(semanticModelEntity.comment())
+ .withSemanticModelDefinition(definitionJson)
+ .withProperties(propertiesJson)
+ .withAuditInfo(
+
JsonUtils.anyFieldMapper().writeValueAsString(semanticModelEntity.auditInfo()))
+ .withDeletedAt(DEFAULT_DELETED_AT)
+ .build();
+ } catch (JsonProcessingException e) {
+ throw new RuntimeException("Failed to serialize Semantic Model JSON", e);
+ }
+ }
+
+ /**
+ * Builds a Semantic Model identity persistent object and the requested
complete snapshot.
+ *
+ * @param semanticModelEntity The Semantic Model entity.
+ * @param builder The identity persistent-object builder.
+ * @param version The version to allocate.
+ * @return The built persistent object.
+ */
+ public static SemanticModelPO buildSemanticModelPO(
+ SemanticModelEntity semanticModelEntity, SemanticModelPOBuilder builder,
Integer version) {
+ try {
+ NamespacedEntityId namespacedEntityId =
+ EntityIdService.getEntityIds(
+ NameIdentifier.of(semanticModelEntity.namespace().levels()),
+ Entity.EntityType.SCHEMA);
+ SemanticModelVersionInfoPO versionPO =
+ initializeSemanticModelVersionInfoPO(semanticModelEntity,
namespacedEntityId, version);
+ return builder
+ .withSemanticModelId(semanticModelEntity.id())
+ .withSemanticModelName(semanticModelEntity.name())
+ .withMetalakeId(namespacedEntityId.namespaceIds()[0])
+ .withCatalogId(namespacedEntityId.namespaceIds()[1])
+ .withSchemaId(namespacedEntityId.entityId())
+ .withAuditInfo(
+
JsonUtils.anyFieldMapper().writeValueAsString(semanticModelEntity.auditInfo()))
+ .withCurrentVersion(version)
+ .withLastVersion(version)
+ .withSemanticModelVersionInfoPO(versionPO)
+ .withDeletedAt(DEFAULT_DELETED_AT)
+ .build();
+ } catch (JsonProcessingException e) {
+ throw new RuntimeException("Failed to serialize Semantic Model audit
info", e);
+ }
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/po/SemanticModelVersionInfoPO.java
b/core/src/main/java/org/apache/gravitino/storage/relational/po/SemanticModelVersionInfoPO.java
new file mode 100644
index 0000000000..03a5491c2a
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/po/SemanticModelVersionInfoPO.java
@@ -0,0 +1,91 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.po;
+
+import com.google.common.base.Preconditions;
+import lombok.EqualsAndHashCode;
+import lombok.Getter;
+import lombok.ToString;
+import lombok.experimental.Accessors;
+import org.apache.commons.lang3.StringUtils;
+
+/** The persistent object for a complete Semantic Model version snapshot. */
+@EqualsAndHashCode
+@Getter
+@ToString
+@Accessors(fluent = true)
+public class SemanticModelVersionInfoPO {
+
+ private Long id;
+ private Long metalakeId;
+ private Long catalogId;
+ private Long schemaId;
+ private Long semanticModelId;
+ private Integer version;
+ private String semanticModelName;
+ private String semanticModelComment;
+ private String semanticModelDefinition;
+ private String properties;
+ private String auditInfo;
+ private Long deletedAt;
+
+ /** Creates an empty persistent object for MyBatis. */
+ public SemanticModelVersionInfoPO() {}
+
+ @lombok.Builder(setterPrefix = "with")
+ private SemanticModelVersionInfoPO(
+ Long id,
+ Long metalakeId,
+ Long catalogId,
+ Long schemaId,
+ Long semanticModelId,
+ Integer version,
+ String semanticModelName,
+ String semanticModelComment,
+ String semanticModelDefinition,
+ String properties,
+ String auditInfo,
+ Long deletedAt) {
+ Preconditions.checkArgument(metalakeId != null, "Metalake id is required");
+ Preconditions.checkArgument(catalogId != null, "Catalog id is required");
+ Preconditions.checkArgument(schemaId != null, "Schema id is required");
+ Preconditions.checkArgument(semanticModelId != null, "Semantic Model id is
required");
+ Preconditions.checkArgument(version != null, "Semantic Model version is
required");
+ Preconditions.checkArgument(
+ StringUtils.isNotBlank(semanticModelName), "Semantic Model name cannot
be empty");
+ Preconditions.checkArgument(
+ StringUtils.isNotBlank(semanticModelDefinition),
+ "Semantic Model definition cannot be empty");
+ Preconditions.checkArgument(StringUtils.isNotBlank(auditInfo), "Audit info
cannot be empty");
+ Preconditions.checkArgument(deletedAt != null, "Deleted at is required");
+
+ this.id = id;
+ this.metalakeId = metalakeId;
+ this.catalogId = catalogId;
+ this.schemaId = schemaId;
+ this.semanticModelId = semanticModelId;
+ this.version = version;
+ this.semanticModelName = semanticModelName;
+ this.semanticModelComment = semanticModelComment;
+ this.semanticModelDefinition = semanticModelDefinition;
+ this.properties = properties;
+ this.auditInfo = auditInfo;
+ this.deletedAt = deletedAt;
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/SemanticModelMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SemanticModelMetaService.java
new file mode 100644
index 0000000000..acc13fbc07
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SemanticModelMetaService.java
@@ -0,0 +1,242 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.service;
+
+import static
org.apache.gravitino.metrics.source.MetricsSource.GRAVITINO_RELATIONAL_STORE_METRIC_NAME;
+import static
org.apache.gravitino.storage.relational.po.SemanticModelPO.fromSemanticModelPO;
+import static
org.apache.gravitino.storage.relational.po.SemanticModelPO.initializeSemanticModelPO;
+
+import com.google.common.base.Preconditions;
+import java.io.IOException;
+import java.util.Locale;
+import java.util.Objects;
+import java.util.concurrent.atomic.AtomicReference;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.meta.SemanticModelEntity;
+import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper;
+import
org.apache.gravitino.storage.relational.mapper.SemanticModelVersionInfoMapper;
+import org.apache.gravitino.storage.relational.po.SemanticModelPO;
+import org.apache.gravitino.storage.relational.po.SemanticModelVersionInfoPO;
+import org.apache.gravitino.storage.relational.utils.ExceptionUtils;
+import org.apache.gravitino.storage.relational.utils.SessionUtils;
+import org.apache.gravitino.utils.NameIdentifierUtil;
+
+/** Provides relational create and load operations for Semantic Model
metadata. */
+public class SemanticModelMetaService {
+
+ private static final SemanticModelMetaService INSTANCE = new
SemanticModelMetaService();
+
+ private final BasePOStorageOps<SemanticModelPO, SemanticModelMetaMapper> ops;
+
+ /** Returns the singleton Semantic Model metadata service. */
+ public static SemanticModelMetaService getInstance() {
+ return INSTANCE;
+ }
+
+ private SemanticModelMetaService() {
+ this.ops = new HierarchicalConversionPOStorageOps<>(new
SemanticModelPOStorageOps());
+ }
+
+ /**
+ * Resolves a Semantic Model stable ID by schema ID and name.
+ *
+ * @param schemaId The parent schema ID.
+ * @param semanticModelName The Semantic Model name.
+ * @return The stable Semantic Model ID.
+ * @throws NoSuchEntityException If the Semantic Model does not exist.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "getSemanticModelIdBySchemaIdAndName")
+ public Long getSemanticModelIdBySchemaIdAndName(long schemaId, String
semanticModelName) {
+ Long semanticModelId =
+ SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper -> mapper.selectSemanticModelIdBySchemaIdAndName(schemaId,
semanticModelName));
+ if (semanticModelId == null) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.SEMANTIC_MODEL.name().toLowerCase(Locale.ROOT),
+ semanticModelName);
+ }
+ return semanticModelId;
+ }
+
+ /**
+ * Loads the current Semantic Model by identifier.
+ *
+ * @param identifier The Semantic Model identifier.
+ * @return The current Semantic Model entity.
+ * @throws NoSuchEntityException If the Semantic Model does not exist.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "getSemanticModelByIdentifier")
+ public SemanticModelEntity getSemanticModelByIdentifier(NameIdentifier
identifier) {
+ SemanticModelPO semanticModelPO =
getSemanticModelPOByIdentifier(identifier);
+ return fromSemanticModelPO(semanticModelPO, identifier.namespace());
+ }
+
+ /**
+ * Inserts or overwrites a Semantic Model identity and its snapshot
atomically.
+ *
+ * @param semanticModelEntity The Semantic Model entity.
+ * @param overwrite Whether to overwrite rows for the same stable ID or
natural key.
+ * @throws IOException If relational persistence fails.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "insertSemanticModel")
+ public void insertSemanticModel(SemanticModelEntity semanticModelEntity,
boolean overwrite)
+ throws IOException {
+
NameIdentifierUtil.checkSemanticModel(semanticModelEntity.nameIdentifier());
+ try {
+ SemanticModelPO po =
+ initializeSemanticModelPO(semanticModelEntity,
SemanticModelPO.builder());
+ AtomicReference<SemanticModelPO> persistedPO = new AtomicReference<>(po);
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ semanticModelEntity.nameIdentifier(),
+ po.getSchemaId(),
+ po.getCatalogId(),
+ po.getMetalakeId(),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper -> {
+ SemanticModelPO storedPO =
+ overwrite
+ ?
mapper.selectSemanticModelMetaBySchemaIdAndNameForUpdate(
+ po.getSchemaId(),
po.getSemanticModelName())
+ : null;
+ if (storedPO == null) {
+ // Keep missing-row imports strict. A losing
concurrent insert must not
+ // overwrite a persisted identity with a different
ID.
+ ops.insertPO(mapper, po, false);
+ return;
+ }
+
+ SemanticModelPO replacementPO =
semanticModelForOverwrite(po, storedPO);
+ Integer updated =
mapper.updateSemanticModelMeta(replacementPO, storedPO);
+ if (updated == null || updated != 1) {
+ throw semanticModelWriteFailure(
+ semanticModelEntity.nameIdentifier(), storedPO);
+ }
+ persistedPO.set(replacementPO);
+ }),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
+ mapper.insertSemanticModelVersionInfo(
+
persistedPO.get().getSemanticModelVersionInfoPO())));
+ } catch (RuntimeException re) {
+ ExceptionUtils.checkSQLException(
+ re, Entity.EntityType.SEMANTIC_MODEL,
semanticModelEntity.nameIdentifier().toString());
+ throw re;
+ }
+ }
+
+ /** Returns the persistent-object operations used by this service. */
+ public BasePOStorageOps<SemanticModelPO, SemanticModelMetaMapper> ops() {
+ return ops;
+ }
+
+ private static SemanticModelPO semanticModelForOverwrite(
+ SemanticModelPO source, SemanticModelPO persistedPO) {
+ int previousVersion = Math.max(persistedPO.getCurrentVersion(),
persistedPO.getLastVersion());
+ Preconditions.checkState(
+ previousVersion < Integer.MAX_VALUE,
+ "Semantic Model %s has exhausted the version range",
+ persistedPO.getSemanticModelId());
+ int nextVersion = previousVersion + 1;
+ return SemanticModelPO.builder()
+ .withSemanticModelId(persistedPO.getSemanticModelId())
+ .withSemanticModelName(source.getSemanticModelName())
+ .withMetalakeId(persistedPO.getMetalakeId())
+ .withCatalogId(persistedPO.getCatalogId())
+ .withSchemaId(persistedPO.getSchemaId())
+ .withAuditInfo(source.getAuditInfo())
+ .withCurrentVersion(nextVersion)
+ .withLastVersion(nextVersion)
+ .withDeletedAt(source.getDeletedAt())
+ .withSemanticModelVersionInfoPO(
+ versionInfoForOverwrite(
+ source.getSemanticModelVersionInfoPO(), persistedPO,
nextVersion))
+ .build();
+ }
+
+ private static SemanticModelVersionInfoPO versionInfoForOverwrite(
+ SemanticModelVersionInfoPO source, SemanticModelPO persistedPO, int
nextVersion) {
+ return SemanticModelVersionInfoPO.builder()
+ .withMetalakeId(persistedPO.getMetalakeId())
+ .withCatalogId(persistedPO.getCatalogId())
+ .withSchemaId(persistedPO.getSchemaId())
+ .withSemanticModelId(persistedPO.getSemanticModelId())
+ .withVersion(nextVersion)
+ .withSemanticModelName(source.semanticModelName())
+ .withSemanticModelComment(source.semanticModelComment())
+ .withSemanticModelDefinition(source.semanticModelDefinition())
+ .withProperties(source.properties())
+ .withAuditInfo(source.auditInfo())
+ .withDeletedAt(source.deletedAt())
+ .build();
+ }
+
+ private SemanticModelPO getSemanticModelPOByIdentifier(NameIdentifier
identifier) {
+ NameIdentifierUtil.checkSemanticModel(identifier);
+ SemanticModelPO semanticModelPO =
+ SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
+ POStorageReadRouting.getPO(
+ mapper, identifier, ops,
Entity.EntityType.SEMANTIC_MODEL));
+ if (semanticModelPO == null) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.SEMANTIC_MODEL.name().toLowerCase(Locale.ROOT),
+ identifier.name());
+ }
+ return semanticModelPO;
+ }
+
+ private RuntimeException semanticModelWriteFailure(
+ NameIdentifier identifier, SemanticModelPO observedSemanticModelPO) {
+ return OccWriteSupport.writeFailure(
+ identifier,
+ Entity.EntityType.SEMANTIC_MODEL,
+ () ->
+ SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
+ mapper.selectSemanticModelMetaByIdForUpdate(
+ observedSemanticModelPO.getSemanticModelId())),
+ null,
+ current ->
+ Objects.equals(
+ current.getSemanticModelName(),
observedSemanticModelPO.getSemanticModelName())
+ && Objects.equals(current.getSchemaId(),
observedSemanticModelPO.getSchemaId())
+ && Objects.equals(current.getCatalogId(),
observedSemanticModelPO.getCatalogId())
+ && Objects.equals(
+ current.getMetalakeId(),
observedSemanticModelPO.getMetalakeId()));
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/SemanticModelPOStorageOps.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SemanticModelPOStorageOps.java
new file mode 100644
index 0000000000..a2877328f7
--- /dev/null
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SemanticModelPOStorageOps.java
@@ -0,0 +1,87 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational.service;
+
+import com.google.common.base.Preconditions;
+import java.util.Locale;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.Namespace;
+import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper;
+import org.apache.gravitino.storage.relational.po.SemanticModelPO;
+
+/** Provides relational persistent-object operations required to create and
load Semantic Models. */
+public class SemanticModelPOStorageOps
+ extends BasePOStorageOps<SemanticModelPO, SemanticModelMetaMapper> {
+
+ /** Creates Semantic Model persistent-object operations. */
+ public SemanticModelPOStorageOps() {}
+
+ /**
+ * {@inheritDoc}
+ *
+ * <p>Overwrite is handled by {@link SemanticModelMetaService} through a
locked, version-checked
+ * update. A database-specific upsert would bypass the persisted identity
and version resolution.
+ */
+ @Override
+ public void insertPO(
+ SemanticModelMetaMapper mapper, SemanticModelPO semanticModelPO, boolean
overwrite) {
+ Preconditions.checkArgument(
+ !overwrite,
+ "Semantic Model overwrite is handled by SemanticModelMetaService, not
by an upsert");
+ mapper.insertSemanticModelMeta(semanticModelPO);
+ }
+
+ @Override
+ public SemanticModelPO getPO(
+ SemanticModelMetaMapper mapper, Long parentId, String semanticModelName)
{
+ return mapper.selectSemanticModelMetaBySchemaIdAndName(parentId,
semanticModelName);
+ }
+
+ @Override
+ public SemanticModelPO getPOByFullName(
+ SemanticModelMetaMapper mapper, NameIdentifier identifier) {
+ Namespace namespace = identifier.namespace();
+ SemanticModelPO po =
+ mapper.selectSemanticModelByFullQualifiedName(
+ namespace.level(0), namespace.level(1), namespace.level(2),
identifier.name());
+ if (po == null) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.CATALOG.name().toLowerCase(Locale.ROOT),
+ namespace.level(1));
+ }
+ if (po.getSchemaId() == null) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.SCHEMA.name().toLowerCase(Locale.ROOT),
+ namespace.level(2));
+ }
+ if (po.getSemanticModelId() == null) {
+ return null;
+ }
+ return po;
+ }
+
+ @Override
+ public boolean supportsParentIdRelationalRead() {
+ return true;
+ }
+}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/TestSemanticModelJDBCBackend.java
b/core/src/test/java/org/apache/gravitino/storage/relational/TestSemanticModelJDBCBackend.java
new file mode 100644
index 0000000000..566c51fea3
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/TestSemanticModelJDBCBackend.java
@@ -0,0 +1,568 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.gravitino.storage.relational;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import com.google.common.collect.ImmutableMap;
+import java.io.IOException;
+import java.math.BigDecimal;
+import java.math.BigInteger;
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.EntityAlreadyExistsException;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.Namespace;
+import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.meta.NamespacedEntityId;
+import org.apache.gravitino.meta.SemanticModelEntity;
+import org.apache.gravitino.semantic.AIContext;
+import org.apache.gravitino.semantic.AIContextObject;
+import org.apache.gravitino.semantic.CustomExtension;
+import org.apache.gravitino.semantic.DataType;
+import org.apache.gravitino.semantic.Dataset;
+import org.apache.gravitino.semantic.DialectExpression;
+import org.apache.gravitino.semantic.Dimension;
+import org.apache.gravitino.semantic.Expression;
+import org.apache.gravitino.semantic.Field;
+import org.apache.gravitino.semantic.Metric;
+import org.apache.gravitino.semantic.Relationship;
+import org.apache.gravitino.semantic.SemanticModelDefinition;
+import org.apache.gravitino.storage.RandomIdGenerator;
+import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
+import org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper;
+import
org.apache.gravitino.storage.relational.mapper.SemanticModelVersionInfoMapper;
+import org.apache.gravitino.storage.relational.po.SemanticModelPO;
+import org.apache.gravitino.storage.relational.service.POStorageReadRouting;
+import
org.apache.gravitino.storage.relational.service.SemanticModelMetaService;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.gravitino.storage.relational.utils.SessionUtils;
+import org.apache.gravitino.utils.NamespaceUtil;
+import org.apache.ibatis.session.SqlSession;
+import org.junit.jupiter.api.TestTemplate;
+
+/** Tests Semantic Model create and load persistence through {@link
JDBCBackend}. */
+public class TestSemanticModelJDBCBackend extends TestJDBCBackend {
+
+ @TestTemplate
+ public void testCreateAndLoadRoundTrip() throws IOException {
+ Namespace namespace = createParents("round_trip");
+ SemanticModelEntity absent =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "absent_optional_model",
+ false,
+ ImmutableMap.of("domain", "sales"));
+ SemanticModelEntity empty =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "empty_optional_model",
+ true,
+ ImmutableMap.of());
+
+ backend.insert(absent, false);
+ backend.insert(empty, false);
+
+ SemanticModelEntity loadedAbsent =
+ backend.get(absent.nameIdentifier(), Entity.EntityType.SEMANTIC_MODEL);
+ SemanticModelEntity loadedEmpty =
+ backend.get(empty.nameIdentifier(), Entity.EntityType.SEMANTIC_MODEL);
+
+ assertEquals(absent, loadedAbsent);
+ assertEquals(empty, loadedEmpty);
+ assertNull(loadedAbsent.definition().relationships());
+ assertNull(loadedAbsent.definition().metrics());
+ assertEquals(0, loadedEmpty.definition().relationships().length);
+ assertEquals(0, loadedEmpty.definition().metrics().length);
+ assertTrue(loadedEmpty.properties().isEmpty());
+ assertEquals(
+ new BigDecimal("1.50"),
+
loadedAbsent.definition().aiContext().object().additionalProperties().get("threshold"));
+ assertNull(loadedAbsent.definition().aiContext().object().examples());
+ assertEquals(0,
loadedAbsent.definition().aiContext().object().synonyms().length);
+ }
+
+ @TestTemplate
+ public void testDuplicateAndMissingEntities() throws IOException {
+ Namespace namespace = createParents("duplicate");
+ SemanticModelEntity original =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "sales_model",
+ false,
+ ImmutableMap.of("owner", "analytics"));
+ backend.insert(original, false);
+
+ SemanticModelEntity duplicate =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ original.name(),
+ true,
+ ImmutableMap.of());
+ assertThrows(EntityAlreadyExistsException.class, () ->
backend.insert(duplicate, false));
+ assertEquals(
+ original, backend.get(original.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+
+ NameIdentifier missingModel = NameIdentifier.of(namespace,
"missing_model");
+ assertThrows(
+ NoSuchEntityException.class,
+ () -> backend.get(missingModel, Entity.EntityType.SEMANTIC_MODEL));
+
+ String suffix = Long.toUnsignedString(RandomIdGenerator.INSTANCE.nextId());
+ String metalakeName = "missing_parent_metalake_" + suffix;
+ String catalogName = "missing_parent_catalog_" + suffix;
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+ Namespace missingParentNamespace =
+ NamespaceUtil.ofSemanticModel(metalakeName, catalogName,
"missing_schema");
+ SemanticModelEntity missingParent =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ missingParentNamespace,
+ "orphan_model",
+ false,
+ ImmutableMap.of());
+
+ assertThrows(NoSuchEntityException.class, () ->
backend.insert(missingParent, false));
+ assertEquals(0, countRows(SemanticModelMetaMapper.TABLE_NAME,
missingParent.id()));
+ assertEquals(0, countRows(SemanticModelVersionInfoMapper.TABLE_NAME,
missingParent.id()));
+ }
+
+ @TestTemplate
+ public void testCreateRollsBackIdentityWhenSnapshotInsertFails() throws
IOException {
+ Namespace namespace = createParents("transaction");
+ SemanticModelEntity semanticModel =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "transaction_model",
+ false,
+ ImmutableMap.of("domain", "finance"));
+ SemanticModelPO po =
+ SemanticModelPO.initializeSemanticModelPO(semanticModel,
SemanticModelPO.builder());
+ SessionUtils.doWithCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
mapper.insertSemanticModelVersionInfo(po.getSemanticModelVersionInfoPO()));
+
+ assertThrows(EntityAlreadyExistsException.class, () ->
backend.insert(semanticModel, false));
+ assertEquals(0, countRows(SemanticModelMetaMapper.TABLE_NAME,
semanticModel.id()));
+ assertEquals(1, countRows(SemanticModelVersionInfoMapper.TABLE_NAME,
semanticModel.id()));
+ assertEquals(0, countEntityChanges());
+ }
+
+ @TestTemplate
+ public void testOverwriteAdvancesVersionAndRetainsPreviousSnapshot() throws
IOException {
+ Namespace namespace = createParents("overwrite_version");
+ SemanticModelEntity original =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "versioned_model",
+ false,
+ ImmutableMap.of("revision", "one"));
+ SemanticModelEntity replacement =
+ semanticModel(
+ original.id(), namespace, original.name(), true,
ImmutableMap.of("revision", "two"));
+
+ backend.insert(original, false);
+ backend.insert(replacement, true);
+
+ assertEquals(
+ replacement, backend.get(original.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ SemanticModelPO persisted = getSemanticModelPO(original.id());
+ assertEquals(2, persisted.getCurrentVersion());
+ assertEquals(2, persisted.getLastVersion());
+ assertEquals(List.of(1, 2), activeSnapshotVersions(original.id()));
+ assertEquals(0, countEntityChanges());
+ }
+
+ @TestTemplate
+ public void testNaturalKeyOverwriteUsesPersistedSemanticModelId() throws
IOException {
+ Namespace namespace = createParents("natural_key_overwrite");
+ SemanticModelEntity original =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "natural_key_model",
+ false,
+ ImmutableMap.of("revision", "one"));
+ SemanticModelEntity replacement =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ original.name(),
+ true,
+ ImmutableMap.of("revision", "two"));
+ SemanticModelEntity expected =
+ semanticModel(
+ original.id(), namespace, original.name(), true,
ImmutableMap.of("revision", "two"));
+
+ backend.insert(original, false);
+ backend.insert(replacement, true);
+
+ assertEquals(
+ expected, backend.get(original.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ SemanticModelPO persisted = getSemanticModelPO(original.id());
+ assertEquals(2, persisted.getCurrentVersion());
+ assertEquals(2, persisted.getLastVersion());
+ assertEquals(List.of(1, 2), activeSnapshotVersions(original.id()));
+ assertEquals(0, countRows(SemanticModelMetaMapper.TABLE_NAME,
replacement.id()));
+ assertEquals(0, countRows(SemanticModelVersionInfoMapper.TABLE_NAME,
replacement.id()));
+ }
+
+ @TestTemplate
+ public void testOverwriteInsertCreatesMissingSemanticModel() throws
IOException {
+ Namespace namespace = createParents("overwrite_insert");
+ SemanticModelEntity imported =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "imported_model",
+ false,
+ ImmutableMap.of("source", "external"));
+
+ backend.insert(imported, true);
+
+ assertEquals(
+ imported, backend.get(imported.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertEquals(1, countRows(SemanticModelMetaMapper.TABLE_NAME,
imported.id()));
+ assertEquals(1, countRows(SemanticModelVersionInfoMapper.TABLE_NAME,
imported.id()));
+ assertEquals(List.of(1), activeSnapshotVersions(imported.id()));
+ }
+
+ @TestTemplate
+ public void testSemanticModelReadRoutesAndEntityIdResolver() throws
IOException {
+ Namespace namespace = createParents("read_routes");
+ SemanticModelEntity semanticModel =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "routed_model",
+ false,
+ ImmutableMap.of("domain", "sales"));
+ backend.insert(semanticModel, false);
+
+ SemanticModelPO persisted = getSemanticModelPO(semanticModel.id());
+ SemanticModelPO lockedById =
+ SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
mapper.selectSemanticModelMetaByIdForUpdate(semanticModel.id()));
+ assertNotNull(lockedById);
+ assertEquals(persisted.getSemanticModelId(),
lockedById.getSemanticModelId());
+ assertEquals(persisted.getSemanticModelName(),
lockedById.getSemanticModelName());
+ assertEquals(persisted.getMetalakeId(), lockedById.getMetalakeId());
+ assertEquals(persisted.getCatalogId(), lockedById.getCatalogId());
+ assertEquals(persisted.getSchemaId(), lockedById.getSchemaId());
+ assertEquals(persisted.getAuditInfo(), lockedById.getAuditInfo());
+ assertEquals(persisted.getCurrentVersion(),
lockedById.getCurrentVersion());
+ assertEquals(persisted.getLastVersion(), lockedById.getLastVersion());
+ assertEquals(persisted.getDeletedAt(), lockedById.getDeletedAt());
+ assertNull(lockedById.getSemanticModelVersionInfoPO());
+
+ SemanticModelPO byParentId =
readSemanticModelPO(semanticModel.nameIdentifier(), true);
+ SemanticModelPO byFullName =
readSemanticModelPO(semanticModel.nameIdentifier(), false);
+ assertEquals(semanticModel.id(), byParentId.getSemanticModelId());
+ assertEquals(semanticModel.id(), byFullName.getSemanticModelId());
+
+ RelationalEntityStoreIdResolver resolver = new
RelationalEntityStoreIdResolver();
+ NamespacedEntityId resolved =
+ resolver.getEntityIds(semanticModel.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL);
+ assertEquals(semanticModel.id().longValue(), resolved.entityId());
+
+ String metalake = namespace.level(0);
+ String catalog = namespace.level(1);
+ String schema = namespace.level(2);
+ List<NameIdentifier> missingParents =
+ List.of(
+ NameIdentifier.of(
+ NamespaceUtil.ofSemanticModel("missing_metalake", catalog,
schema),
+ semanticModel.name()),
+ NameIdentifier.of(
+ NamespaceUtil.ofSemanticModel(metalake, "missing_catalog",
schema),
+ semanticModel.name()),
+ NameIdentifier.of(
+ NamespaceUtil.ofSemanticModel(metalake, catalog,
"missing_schema"),
+ semanticModel.name()));
+ for (NameIdentifier missing : missingParents) {
+ assertThrows(NoSuchEntityException.class, () ->
readSemanticModelPO(missing, true));
+ assertThrows(NoSuchEntityException.class, () ->
readSemanticModelPO(missing, false));
+ assertThrows(
+ NoSuchEntityException.class,
+ () -> resolver.getEntityIds(missing,
Entity.EntityType.SEMANTIC_MODEL));
+ }
+
+ NameIdentifier missingModel = NameIdentifier.of(namespace,
"missing_model");
+ assertNull(readSemanticModelPO(missingModel, true));
+ assertNull(readSemanticModelPO(missingModel, false));
+ assertThrows(
+ NoSuchEntityException.class,
+ () -> resolver.getEntityIds(missingModel,
Entity.EntityType.SEMANTIC_MODEL));
+ }
+
+ @TestTemplate
+ public void testConcurrentCreatesAndOverwriteRacesRemainAtomic() throws
Exception {
+ Namespace namespace = createParents("concurrent");
+ SemanticModelEntity first =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "concurrent_model",
+ false,
+ ImmutableMap.of("candidate", "first"));
+ SemanticModelEntity second =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ first.name(),
+ true,
+ ImmutableMap.of("candidate", "second"));
+
+ List<Throwable> createResults = insertConcurrently(first, false, second,
false);
+ assertEquals(1, createResults.stream().filter(Objects::isNull).count());
+ Throwable createFailure =
+
createResults.stream().filter(Objects::nonNull).findFirst().orElseThrow();
+ assertTrue(createFailure instanceof EntityAlreadyExistsException);
+ assertEquals(
+ 1,
+ countRows(SemanticModelMetaMapper.TABLE_NAME, first.id())
+ + countRows(SemanticModelMetaMapper.TABLE_NAME, second.id()));
+ assertEquals(
+ 1,
+ countRows(SemanticModelVersionInfoMapper.TABLE_NAME, first.id())
+ + countRows(SemanticModelVersionInfoMapper.TABLE_NAME,
second.id()));
+
+ SemanticModelEntity created =
+ backend.get(first.nameIdentifier(), Entity.EntityType.SEMANTIC_MODEL);
+ SemanticModelEntity overwriteOne =
+ semanticModel(
+ created.id(), namespace, created.name(), false,
ImmutableMap.of("winner", "one"));
+ SemanticModelEntity overwriteTwo =
+ semanticModel(
+ created.id(), namespace, created.name(), true,
ImmutableMap.of("winner", "two"));
+
+ List<Throwable> overwriteResults = insertConcurrently(overwriteOne, true,
overwriteTwo, true);
+ assertTrue(overwriteResults.stream().allMatch(Objects::isNull));
+
+ SemanticModelPO persisted = getSemanticModelPO(created.id());
+ assertEquals(3, persisted.getCurrentVersion());
+ assertEquals(3, persisted.getLastVersion());
+ assertEquals(List.of(1, 2, 3), activeSnapshotVersions(created.id()));
+ SemanticModelEntity winner =
+ backend.get(created.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL);
+ assertTrue(winner.equals(overwriteOne) || winner.equals(overwriteTwo));
+ assertEquals(0, countEntityChanges());
+ }
+
+ private Namespace createParents(String prefix) throws IOException {
+ String suffix = Long.toUnsignedString(RandomIdGenerator.INSTANCE.nextId());
+ String metalakeName = prefix + "_metalake_" + suffix;
+ String catalogName = prefix + "_catalog_" + suffix;
+ String schemaName = prefix + "_schema_" + suffix;
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+ createAndInsertSchema(metalakeName, catalogName, schemaName);
+ return NamespaceUtil.ofSemanticModel(metalakeName, catalogName,
schemaName);
+ }
+
+ private SemanticModelEntity semanticModel(
+ Long id,
+ Namespace namespace,
+ String name,
+ boolean explicitEmpty,
+ Map<String, String> properties) {
+ AIContextObject context =
+ AIContextObject.builder()
+ .withInstructions("Use certified sales definitions")
+ .withSynonyms(new String[0])
+ .withAdditionalProperties(
+ Map.of("threshold", new BigDecimal("1.50"), "nested",
List.of(new BigInteger("3"))))
+ .build();
+ Field field =
+ Field.builder()
+ .withName("ordered_at")
+ .withExpression(
+ Expression.builder()
+ .withDialects(
+ new DialectExpression[] {
+ DialectExpression.builder()
+ .withDialect("ansi")
+ .withExpression("ordered_at")
+ .build()
+ })
+ .build())
+ .withDimension(Dimension.builder().withIsTime(true).build())
+ .withDatatype(DataType.DATE_TIME_TZ)
+ .build();
+ Dataset dataset =
+ Dataset.builder()
+ .withName("orders")
+ .withSource(NameIdentifier.of("sales", "mart", "orders"))
+ .withPrimaryKey(new String[0])
+ .withFields(new Field[] {field})
+ .build();
+ SemanticModelDefinition.Builder definitionBuilder =
+ SemanticModelDefinition.builder()
+ .withAIContext(AIContext.of(context))
+ .withDatasets(new Dataset[] {dataset});
+ if (explicitEmpty) {
+ definitionBuilder
+ .withRelationships(new Relationship[0])
+ .withMetrics(new Metric[0])
+ .withCustomExtensions(new CustomExtension[0]);
+ }
+
+ return SemanticModelEntity.builder()
+ .withId(id)
+ .withName(name)
+ .withNamespace(namespace)
+ .withComment(explicitEmpty ? null : "Governed sales definitions")
+ .withDefinition(definitionBuilder.build())
+ .withProperties(properties)
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ }
+
+ private int countRows(String tableName, Long semanticModelId) {
+ if (!SemanticModelMetaMapper.TABLE_NAME.equals(tableName)
+ && !SemanticModelVersionInfoMapper.TABLE_NAME.equals(tableName)) {
+ throw new IllegalArgumentException("Unsupported Semantic Model table: "
+ tableName);
+ }
+ String sql = String.format("SELECT count(*) FROM %s WHERE
semantic_model_id = ?", tableName);
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ PreparedStatement statement = connection.prepareStatement(sql)) {
+ statement.setLong(1, semanticModelId);
+ try (ResultSet resultSet = statement.executeQuery()) {
+ assertTrue(resultSet.next());
+ return resultSet.getInt(1);
+ }
+ } catch (SQLException e) {
+ throw new RuntimeException("Failed to count Semantic Model rows", e);
+ }
+ }
+
+ private SemanticModelPO getSemanticModelPO(Long semanticModelId) {
+ SemanticModelPO po =
+ SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper -> mapper.selectSemanticModelMetaById(semanticModelId));
+ assertNotNull(po);
+ return po;
+ }
+
+ private SemanticModelPO readSemanticModelPO(NameIdentifier identifier,
boolean cacheEnabled) {
+ return SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
+ POStorageReadRouting.getPO(
+ mapper,
+ identifier,
+ SemanticModelMetaService.getInstance().ops(),
+ Entity.EntityType.SEMANTIC_MODEL,
+ cacheEnabled));
+ }
+
+ private List<Integer> activeSnapshotVersions(Long semanticModelId) {
+ String sql =
+ String.format(
+ "SELECT version FROM %s WHERE semantic_model_id = ? AND deleted_at
= 0 ORDER BY version",
+ SemanticModelVersionInfoMapper.TABLE_NAME);
+ List<Integer> versions = new ArrayList<>();
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ PreparedStatement statement = connection.prepareStatement(sql)) {
+ statement.setLong(1, semanticModelId);
+ try (ResultSet resultSet = statement.executeQuery()) {
+ while (resultSet.next()) {
+ versions.add(resultSet.getInt(1));
+ }
+ }
+ return versions;
+ } catch (SQLException e) {
+ throw new RuntimeException("Failed to list Semantic Model versions", e);
+ }
+ }
+
+ private int countEntityChanges() {
+ return SessionUtils.getWithoutCommit(
+ EntityChangeLogMapper.class,
+ mapper ->
+ Math.toIntExact(
+ mapper.selectEntityChanges(0, 100).stream()
+ .filter(
+ record ->
+
Entity.EntityType.SEMANTIC_MODEL.name().equals(record.getEntityType()))
+ .count()));
+ }
+
+ private List<Throwable> insertConcurrently(
+ SemanticModelEntity first,
+ boolean overwriteFirst,
+ SemanticModelEntity second,
+ boolean overwriteSecond)
+ throws Exception {
+ CountDownLatch start = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ try {
+ Future<Throwable> firstResult =
+ executor.submit(() -> insertAfterStart(first, overwriteFirst,
start));
+ Future<Throwable> secondResult =
+ executor.submit(() -> insertAfterStart(second, overwriteSecond,
start));
+ start.countDown();
+ return Arrays.asList(
+ firstResult.get(30, TimeUnit.SECONDS), secondResult.get(30,
TimeUnit.SECONDS));
+ } finally {
+ executor.shutdownNow();
+ }
+ }
+
+ private Throwable insertAfterStart(
+ SemanticModelEntity semanticModel, boolean overwrite, CountDownLatch
start) {
+ try {
+ assertTrue(start.await(30, TimeUnit.SECONDS));
+ backend.insert(semanticModel, overwrite);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ }
+}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
index 60a0dded70..cda36c09f7 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
@@ -59,12 +59,15 @@ import org.apache.gravitino.meta.FunctionEntity;
import org.apache.gravitino.meta.ModelEntity;
import org.apache.gravitino.meta.ModelVersionEntity;
import org.apache.gravitino.meta.SchemaEntity;
+import org.apache.gravitino.meta.SemanticModelEntity;
import org.apache.gravitino.meta.TableEntity;
import org.apache.gravitino.meta.TagEntity;
import org.apache.gravitino.meta.TopicEntity;
import org.apache.gravitino.meta.ViewEntity;
import org.apache.gravitino.model.ModelVersion;
import org.apache.gravitino.rel.types.Types;
+import org.apache.gravitino.semantic.Dataset;
+import org.apache.gravitino.semantic.SemanticModelDefinition;
import org.apache.gravitino.storage.RandomIdGenerator;
import org.apache.gravitino.storage.relational.RelationalBackend;
import org.apache.gravitino.storage.relational.TestJDBCBackend;
@@ -160,6 +163,37 @@ public class TestSchemaMetaService extends TestJDBCBackend
{
assertSchemaChildActionWaitsForConcurrentDelete(
schema, () -> childCase.write.run(childNamespace));
}
+
+ SchemaEntity semanticModelSchema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ "schema_for_entity_lock_semantic_model",
+ AUDIT_INFO);
+ backend.insert(semanticModelSchema, false);
+ Namespace semanticModelNamespace =
+ Namespace.of(metalakeName, catalogName, semanticModelSchema.name());
+ assertSchemaChildActionWaitsForConcurrentDelete(
+ semanticModelSchema,
+ () ->
+ backend.insert(
+ SemanticModelEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("child_semantic_model")
+ .withNamespace(semanticModelNamespace)
+ .withDefinition(
+ SemanticModelDefinition.builder()
+ .withDatasets(
+ new Dataset[] {
+ Dataset.builder()
+ .withName("child_dataset")
+
.withSource(NameIdentifier.of("source_table"))
+ .build()
+ })
+ .build())
+ .withAuditInfo(AUDIT_INFO)
+ .build(),
+ false));
}
@TestTemplate
diff --git a/design-docs/gravitino-entity-cache-multinode-design.md
b/design-docs/gravitino-entity-cache-multinode-design.md
index 714d8ecb03..949d85530f 100644
--- a/design-docs/gravitino-entity-cache-multinode-design.md
+++ b/design-docs/gravitino-entity-cache-multinode-design.md
@@ -274,7 +274,7 @@ The danger depends on where the real data comes from.
**Group 1 — the data comes from the connector: table, view, topic, and schemas
in external catalogs.** For these, Gravitino does not read the metadata from
its own cache. On every load it asks the underlying system (Hive, Iceberg,
JDBC, Kafka) for the current object, and uses the cached entity only for
Gravitino's own id and audit fields. All nodes ask the same underlying system,
so they all see the same thing. If the object was dropped or renamed, the
connector call fails right away wit [...]
-**Group 2 — the data comes from Gravitino's own store: metalake, catalog,
managed schema, fileset, model, model version, function.** Here the cached
entity *is* the answer, so a stale read really can return an old value. We
looked at each alter in this group.
+**Group 2 — the data comes from Gravitino's own store: metalake, catalog,
managed schema, fileset, model, model version, Semantic Model, function.** Here
the cached entity *is* the answer, so a stale read really can return an old
value. We looked at each alter in this group.
(Views follow the same rule as tables: their definition is read through the
connector on every load, so they sit in Group 1 and are safe.)
@@ -291,22 +291,23 @@ Most alters are safe:
| drop | for Group 1 the
connector call fails; for a fileset the file path is gone and access fails
| yes |
| catalog changes | already handled today
by a separate cross-node signal (`CatalogChangeLogListener`), not by this cache
| yes |
-Five cases are **not** safe — a stale read gives a wrong answer with no error.
They are in the model area, functions, or the metalake on/off flag:
+Six cases are **not** safe — a stale read gives a wrong answer with no error.
They are in the model area, Semantic Models, functions, or the metalake on/off
flag:
| Change | Why
a stale read is wrong, with no error
|
| --------------------------------------------------------------------- |
---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
|
| model version — update / add / remove URI | the
URI points to the model files. An old URI silently loads the **wrong model
files**.
|
| model version — update aliases | an
alias silently points to the **wrong version**.
|
| model — latest version (after a new version is added on another node) | "get
the latest version" silently returns an **old version**.
|
+| Semantic Model — overwrite definition | the
current version selects the definition used for semantic queries. An old
version silently uses the **wrong definition**.
|
| function — add / update / remove implementation or definition | the
implementation is the code the function runs. An old copy silently runs the
**wrong function code**.
|
| metalake — disable | a
node with a stale copy still thinks the metalake is on and **lets operations
run** on a metalake that was turned off. (Turning it back on is safe: a stale
node only throws "in use" by mistake.) |
#### What we do about it
-A stale read is only a problem for a **per-node** cache, and only for the
load-bearing pointers listed above. A **shared** cache (redis) keeps one copy
for the whole cluster, so it has no window at all and can safely cache
everything. So the handling depends on the cache kind (`coherence()`):
+A stale read is only a problem for a **per-node** cache, and only for the
load-bearing content listed above. A **shared** cache (redis) keeps one copy
for the whole cluster, so it has no window at all and can safely cache
everything. So the handling depends on the cache kind (`coherence()`):
-- **Shared cache (redis, `SHARED`): cache everything.** There is no per-node
window, so model, model version, and function are cached like any other entity,
with no extra work. The writing node clears the one shared copy, and every node
sees it at once.
-- **Per-node cache (caffeine, `LOCAL_PER_NODE`): do not cache model, model
version, or function.** Each holds a load-bearing pointer (a version URI, the
latest version, or the function implementation) that would be silently wrong on
another node during the poll window. They are read rarely, so reading them from
the DB every time costs little, and it keeps the rule simple — an entity type
is either in or out, with no special per-read handling. If they ever get hot,
we can revisit.
+- **Shared cache (redis, `SHARED`): cache everything.** There is no per-node
window, so model, model version, Semantic Model, and function are cached like
any other entity, with no extra work. The writing node clears the one shared
copy, and every node sees it at once.
+- **Per-node cache (caffeine, `LOCAL_PER_NODE`): do not cache model, model
version, Semantic Model, or function.** Each holds load-bearing content (a
version URI, the latest version, a Semantic Model definition, or a function
implementation) that would be silently wrong on another node during the poll
window. They are read rarely, so reading them from the DB every time costs
little, and it keeps the rule simple — an entity type is either in or out, with
no special per-read handling. If t [...]
- **Metalake on/off flag — cached like the rest of the metalake.** Disabling
or deleting a metalake is a rare, tenant-level admin action, so we accept the
small window instead of adding special handling: the metalake is cached and
invalidated across nodes through the change log like any other entity, so after
a disable, another node stops allowing operations within one poll interval.
We do **not** need a per-entity version check on the cache: for a point read,
checking the DB version costs the same query as just reading the row, so it
would buy nothing.
@@ -322,10 +323,11 @@ We do **not** need a per-entity version check on the
cache: for a point read, ch
| tag, policy | self-contained; a stale read is only cosmetic (old
comment/property) | cache + change-log
invalidation | cache |
| job | self-contained operational entity; a stale read is
only an old job status | cache + change-log
invalidation | cache |
| model, model version | carries a load-bearing pointer (a version's URI, the
latest version) that would be wrong if stale | **not cached — read from the
DB** (revisit if it gets hot) | cache |
+| Semantic Model | its current version selects the definition used for
semantic queries | **not cached — read from the
DB** (revisit if it gets hot) | cache |
| function | the cached value *is* the code that runs, so a stale
copy would run the wrong code | **not cached — read from the
DB** (revisit if it gets hot) | cache |
| user, group, role | derived fields (`roleNames`, `securableObjects`) need
a reverse lookup | not cached
| not cached |
-After this, everything a per-node cache serves is safe or bounded:
connector-backed entities are safe by construction; self-contained store
entities are only ever cosmetically stale (and a rare metalake disable
self-corrects within one poll interval); model / model version / function are
read from the DB. A shared cache is safe throughout because it has no window.
+After this, everything a per-node cache serves is safe or bounded:
connector-backed entities are safe by construction; self-contained store
entities are only ever cosmetically stale (and a rare metalake disable
self-corrects within one poll interval); model / model version / Semantic Model
/ function are read from the DB. A shared cache is safe throughout because it
has no window.
#### The staleness promise (SLA)