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 52ab3fabfc [#12601] feat(store): Add Semantic Model list, alter, and
drop support (#12603)
52ab3fabfc is described below
commit 52ab3fabfc54db1d717b58a67a7b08d204b52028
Author: mchades <[email protected]>
AuthorDate: Mon Sep 28 15:41:07 2026 +0800
[#12601] feat(store): Add Semantic Model list, alter, and drop support
(#12603)
### What changes were proposed in this pull request?
Add relational persistence for listing, altering, and dropping Semantic
Models, including optimistic locking, cascade handling, hard deletion,
version retention, and cleanup integration.
### Why are the changes needed?
Semantic Models need complete relational lifecycle support beyond create
and load operations.
Fix: #12601
### 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, Java compilation, test compilation, and `git diff
--check` passed.
- `TestSemanticModelJDBCBackend`: 33/33 across H2, MySQL, and
PostgreSQL.
- `TestSemanticModelMetaService`: 21/21 across H2, MySQL, and
PostgreSQL.
- `TestSchemaMetaService`: H2 run passed.
- `TestBaseEntityCache`: 6/6.
---
.../gravitino/storage/relational/JDBCBackend.java | 30 +-
.../relational/mapper/SemanticModelMetaMapper.java | 70 ++-
.../SemanticModelMetaSQLProviderFactory.java | 48 +-
.../mapper/SemanticModelVersionInfoMapper.java | 61 +-
...SemanticModelVersionInfoSQLProviderFactory.java | 56 +-
.../provider/base/SchemaMetaBaseSQLProvider.java | 4 +
.../base/SemanticModelMetaBaseSQLProvider.java | 150 +++--
.../SemanticModelVersionInfoBaseSQLProvider.java | 85 ++-
.../SemanticModelMetaPostgreSQLProvider.java | 56 +-
...SemanticModelVersionInfoPostgreSQLProvider.java | 69 ++-
.../relational/service/CatalogMetaService.java | 12 +-
.../relational/service/MetalakeMetaService.java | 12 +-
.../relational/service/SchemaMetaService.java | 12 +-
.../service/SemanticModelMetaService.java | 285 ++++++++-
.../service/SemanticModelPOStorageOps.java | 36 +-
.../storage/relational/TestJDBCBackend.java | 4 +
.../relational/TestSemanticModelJDBCBackend.java | 174 +++++-
.../base/TestSchemaMetaBaseSQLProvider.java | 4 +-
.../service/TestSemanticModelMetaService.java | 687 +++++++++++++++++++++
19 files changed, 1779 insertions(+), 76 deletions(-)
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 291f00535d..b6f184fa8c 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
@@ -142,6 +142,9 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
return (List<E>)
TableMetaService.getInstance().listTablesByNamespace(namespace);
case VIEW:
return (List<E>)
ViewMetaService.getInstance().listViewsByNamespace(namespace);
+ case SEMANTIC_MODEL:
+ return (List<E>)
+
SemanticModelMetaService.getInstance().listSemanticModelsByNamespace(namespace);
case FILESET:
return (List<E>)
FilesetMetaService.getInstance().listFilesetsByNamespace(namespace);
case TOPIC:
@@ -362,6 +365,19 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
}
}
return views;
+ case SEMANTIC_MODEL:
+ List<E> semanticModels = Lists.newArrayList();
+ for (NameIdentifier identifier : identifiers) {
+ try {
+ semanticModels.add(
+ (E)
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(identifier));
+ } catch (NoSuchEntityException e) {
+ LOG.debug("Skipping missing semantic model during batch get: {}",
identifier.name());
+ }
+ }
+ return semanticModels;
default:
throw new UnsupportedEntityTypeException(
"Unsupported entity type: %s for batch get operation", entityType);
@@ -519,8 +535,9 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
.deleteViewMetasByLegacyTimeline(
legacyTimeline, GARBAGE_COLLECTOR_SINGLE_DELETION_LIMIT);
case SEMANTIC_MODEL:
- // TODO(#12209): Delegate to SemanticModelMetaService when relational
persistence is added.
- return 0;
+ return SemanticModelMetaService.getInstance()
+ .deleteSemanticModelMetasByLegacyTimeline(
+ legacyTimeline, GARBAGE_COLLECTOR_SINGLE_DELETION_LIMIT);
case AUDIT:
return 0;
// TODO: Implement hard delete logic for these entity types.
@@ -562,8 +579,9 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
return 0;
case SEMANTIC_MODEL:
- // TODO: Delegate to SemanticModelMetaService when relational
persistence is added.
- return 0;
+ return SemanticModelMetaService.getInstance()
+ .deleteSemanticModelVersionsByRetentionCount(
+ versionRetentionCount,
GARBAGE_COLLECTOR_SINGLE_DELETION_LIMIT);
case FILESET:
return FilesetMetaService.getInstance()
@@ -1048,6 +1066,8 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
return (E) JobMetaService.getInstance().updateJob(ident, updater);
case VIEW:
return (E) ViewMetaService.getInstance().updateView(ident, updater);
+ case SEMANTIC_MODEL:
+ return (E)
SemanticModelMetaService.getInstance().updateSemanticModel(ident, updater);
default:
throw new UnsupportedEntityTypeException(
"Unsupported entity type: %s for update operation", entityType);
@@ -1091,6 +1111,8 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
return JobMetaService.getInstance().deleteJob(ident);
case VIEW:
return ViewMetaService.getInstance().deleteView(ident);
+ case SEMANTIC_MODEL:
+ return
SemanticModelMetaService.getInstance().deleteSemanticModel(ident);
default:
throw new UnsupportedEntityTypeException(
"Unsupported entity type: %s for delete operation", entityType);
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
index a8690f8230..37c4c0ceab 100644
---
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
@@ -18,8 +18,10 @@
*/
package org.apache.gravitino.storage.relational.mapper;
+import java.util.List;
import org.apache.gravitino.storage.relational.po.SemanticModelPO;
import org.apache.gravitino.storage.relational.po.SemanticModelVersionInfoPO;
+import org.apache.ibatis.annotations.DeleteProvider;
import org.apache.ibatis.annotations.InsertProvider;
import org.apache.ibatis.annotations.One;
import org.apache.ibatis.annotations.Param;
@@ -30,7 +32,7 @@ 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. */
+/** A MyBatis mapper for Semantic Model identity metadata operations. */
public interface SemanticModelMetaMapper {
/** The Semantic Model identity table name. */
@@ -59,14 +61,7 @@ public interface SemanticModelMetaMapper {
@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. */
+ /** Lists current Semantic Model snapshots under a schema ID. */
@Results(
id = "semanticModelPOResultMap",
value = {
@@ -89,6 +84,30 @@ public interface SemanticModelMetaMapper {
+ "version_audit_info,version_deleted_at}",
one = @One(resultMap = "mapToSemanticModelVersionInfoPO"))
})
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "listSemanticModelPOsBySchemaId")
+ List<SemanticModelPO> listSemanticModelPOsBySchemaId(@Param("schemaId") Long
schemaId);
+
+ /** Lists current Semantic Model snapshots under a fully qualified schema
name. */
+ @ResultMap("semanticModelPOResultMap")
+ @SelectProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "listSemanticModelPOsByFullQualifiedName")
+ List<SemanticModelPO> listSemanticModelPOsByFullQualifiedName(
+ @Param("metalakeName") String metalakeName,
+ @Param("catalogName") String catalogName,
+ @Param("schemaName") String schemaName);
+
+ /** 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. */
+ @ResultMap("semanticModelPOResultMap")
@SelectProvider(
type = SemanticModelMetaSQLProviderFactory.class,
method = "selectSemanticModelMetaBySchemaIdAndName")
@@ -140,4 +159,37 @@ public interface SemanticModelMetaMapper {
Integer updateSemanticModelMeta(
@Param("newSemanticModelMeta") SemanticModelPO newSemanticModelPO,
@Param("oldSemanticModelMeta") SemanticModelPO oldSemanticModelPO);
+
+ /** Soft-deletes a Semantic Model identity by stable ID and expected current
version. */
+ @UpdateProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "softDeleteSemanticModelMetasBySemanticModelId")
+ Integer softDeleteSemanticModelMetasBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("currentVersion") Integer currentVersion);
+
+ /** Soft-deletes Semantic Model identities under a metalake. */
+ @UpdateProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "softDeleteSemanticModelMetasByMetalakeId")
+ Integer softDeleteSemanticModelMetasByMetalakeId(@Param("metalakeId") Long
metalakeId);
+
+ /** Soft-deletes Semantic Model identities under a catalog. */
+ @UpdateProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "softDeleteSemanticModelMetasByCatalogId")
+ Integer softDeleteSemanticModelMetasByCatalogId(@Param("catalogId") Long
catalogId);
+
+ /** Soft-deletes Semantic Model identities under schemas. */
+ @UpdateProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "softDeleteSemanticModelMetasBySchemaIds")
+ Integer softDeleteSemanticModelMetasBySchemaIds(@Param("schemaIds")
List<Long> schemaIds);
+
+ /** Permanently deletes soft-deleted Semantic Model identities older than a
timeline. */
+ @DeleteProvider(
+ type = SemanticModelMetaSQLProviderFactory.class,
+ method = "deleteSemanticModelMetasByLegacyTimeline")
+ Integer deleteSemanticModelMetasByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit);
}
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
index ca3e3d5865..b0b7fe4d6a 100644
---
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
@@ -19,6 +19,7 @@
package org.apache.gravitino.storage.relational.mapper;
import com.google.common.collect.ImmutableMap;
+import java.util.List;
import java.util.Map;
import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType;
import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelMetaBaseSQLProvider;
@@ -27,7 +28,7 @@ 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. */
+/** Selects database-specific SQL providers for Semantic Model identity
metadata. */
public class SemanticModelMetaSQLProviderFactory {
private static final Map<JDBCBackendType, SemanticModelMetaBaseSQLProvider>
@@ -51,6 +52,20 @@ public class SemanticModelMetaSQLProviderFactory {
static class SemanticModelMetaH2Provider extends
SemanticModelMetaBaseSQLProvider {}
+ /** Provides SQL for listing Semantic Models by schema ID. */
+ public static String listSemanticModelPOsBySchemaId(@Param("schemaId") Long
schemaId) {
+ return getProvider().listSemanticModelPOsBySchemaId(schemaId);
+ }
+
+ /** Provides SQL for listing Semantic Models by fully qualified schema name.
*/
+ public static String listSemanticModelPOsByFullQualifiedName(
+ @Param("metalakeName") String metalakeName,
+ @Param("catalogName") String catalogName,
+ @Param("schemaName") String schemaName) {
+ return getProvider()
+ .listSemanticModelPOsByFullQualifiedName(metalakeName, catalogName,
schemaName);
+ }
+
/** 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) {
@@ -104,4 +119,35 @@ public class SemanticModelMetaSQLProviderFactory {
@Param("oldSemanticModelMeta") SemanticModelPO oldSemanticModelPO) {
return getProvider().updateSemanticModelMeta(newSemanticModelPO,
oldSemanticModelPO);
}
+
+ /** Provides SQL for soft-deleting a Semantic Model identity by stable ID
and current version. */
+ public static String softDeleteSemanticModelMetasBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("currentVersion") Integer currentVersion) {
+ return getProvider()
+ .softDeleteSemanticModelMetasBySemanticModelId(semanticModelId,
currentVersion);
+ }
+
+ /** Provides SQL for soft-deleting Semantic Model identities by metalake ID.
*/
+ public static String softDeleteSemanticModelMetasByMetalakeId(
+ @Param("metalakeId") Long metalakeId) {
+ return getProvider().softDeleteSemanticModelMetasByMetalakeId(metalakeId);
+ }
+
+ /** Provides SQL for soft-deleting Semantic Model identities by catalog ID.
*/
+ public static String
softDeleteSemanticModelMetasByCatalogId(@Param("catalogId") Long catalogId) {
+ return getProvider().softDeleteSemanticModelMetasByCatalogId(catalogId);
+ }
+
+ /** Provides SQL for soft-deleting Semantic Model identities by schema IDs.
*/
+ public static String softDeleteSemanticModelMetasBySchemaIds(
+ @Param("schemaIds") List<Long> schemaIds) {
+ return getProvider().softDeleteSemanticModelMetasBySchemaIds(schemaIds);
+ }
+
+ /** Provides SQL for permanently deleting old Semantic Model identities. */
+ public static String deleteSemanticModelMetasByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
+ return
getProvider().deleteSemanticModelMetasByLegacyTimeline(legacyTimeline, limit);
+ }
}
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
index c02b2c43be..81437cc584 100644
---
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
@@ -18,11 +18,15 @@
*/
package org.apache.gravitino.storage.relational.mapper;
+import java.util.List;
import org.apache.gravitino.storage.relational.po.SemanticModelVersionInfoPO;
+import org.apache.ibatis.annotations.DeleteProvider;
import org.apache.ibatis.annotations.InsertProvider;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.annotations.SelectProvider;
+import org.apache.ibatis.annotations.UpdateProvider;
-/** A MyBatis mapper for creating Semantic Model version snapshots. */
+/** A MyBatis mapper for Semantic Model version snapshot operations. */
public interface SemanticModelVersionInfoMapper {
/** The Semantic Model version snapshot table name. */
@@ -34,4 +38,59 @@ public interface SemanticModelVersionInfoMapper {
method = "insertSemanticModelVersionInfo")
void insertSemanticModelVersionInfo(
@Param("semanticModelVersionInfo") SemanticModelVersionInfoPO
versionInfoPO);
+
+ /** Selects a Semantic Model version snapshot. */
+ @SelectProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "selectSemanticModelVersionInfoBySemanticModelIdAndVersion")
+ SemanticModelVersionInfoPO
selectSemanticModelVersionInfoBySemanticModelIdAndVersion(
+ @Param("semanticModelId") Long semanticModelId, @Param("version")
Integer version);
+
+ /** Soft-deletes all snapshots for a Semantic Model ID. */
+ @UpdateProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "softDeleteSemanticModelVersionsBySemanticModelId")
+ Integer softDeleteSemanticModelVersionsBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId);
+
+ /** Soft-deletes Semantic Model snapshots under schemas. */
+ @UpdateProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "softDeleteSemanticModelVersionsBySchemaIds")
+ Integer softDeleteSemanticModelVersionsBySchemaIds(@Param("schemaIds")
List<Long> schemaIds);
+
+ /** Soft-deletes Semantic Model snapshots under a catalog. */
+ @UpdateProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "softDeleteSemanticModelVersionsByCatalogId")
+ Integer softDeleteSemanticModelVersionsByCatalogId(@Param("catalogId") Long
catalogId);
+
+ /** Soft-deletes Semantic Model snapshots under a metalake. */
+ @UpdateProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "softDeleteSemanticModelVersionsByMetalakeId")
+ Integer softDeleteSemanticModelVersionsByMetalakeId(@Param("metalakeId")
Long metalakeId);
+
+ /** Permanently deletes soft-deleted snapshots older than a timeline. */
+ @DeleteProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "deleteSemanticModelVersionsByLegacyTimeline")
+ Integer deleteSemanticModelVersionsByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit);
+
+ /** Selects active Semantic Models whose version count exceeds the retention
count. */
+ @SelectProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "selectSemanticModelVersionsByRetentionCount")
+ List<SemanticModelVersionInfoPO> selectSemanticModelVersionsByRetentionCount(
+ @Param("versionRetentionCount") Long versionRetentionCount);
+
+ /** Soft-deletes old snapshots through a per-model retention line. */
+ @UpdateProvider(
+ type = SemanticModelVersionInfoSQLProviderFactory.class,
+ method = "softDeleteSemanticModelVersionsByRetentionLine")
+ Integer softDeleteSemanticModelVersionsByRetentionLine(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("versionRetentionLine") long versionRetentionLine,
+ @Param("limit") int limit);
}
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
index 7e61cfc110..d069d814ab 100644
---
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
@@ -19,6 +19,7 @@
package org.apache.gravitino.storage.relational.mapper;
import com.google.common.collect.ImmutableMap;
+import java.util.List;
import java.util.Map;
import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType;
import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelVersionInfoBaseSQLProvider;
@@ -27,7 +28,7 @@ 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. */
+/** Selects database-specific SQL providers for Semantic Model version
snapshots. */
public class SemanticModelVersionInfoSQLProviderFactory {
private static final Map<JDBCBackendType,
SemanticModelVersionInfoBaseSQLProvider>
@@ -57,4 +58,57 @@ public class SemanticModelVersionInfoSQLProviderFactory {
@Param("semanticModelVersionInfo") SemanticModelVersionInfoPO
versionInfoPO) {
return getProvider().insertSemanticModelVersionInfo(versionInfoPO);
}
+
+ /** Provides SQL for selecting a Semantic Model version snapshot. */
+ public static String
selectSemanticModelVersionInfoBySemanticModelIdAndVersion(
+ @Param("semanticModelId") Long semanticModelId, @Param("version")
Integer version) {
+ return getProvider()
+
.selectSemanticModelVersionInfoBySemanticModelIdAndVersion(semanticModelId,
version);
+ }
+
+ /** Provides SQL for soft-deleting snapshots by Semantic Model ID. */
+ public static String softDeleteSemanticModelVersionsBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId) {
+ return
getProvider().softDeleteSemanticModelVersionsBySemanticModelId(semanticModelId);
+ }
+
+ /** Provides SQL for soft-deleting snapshots by schema IDs. */
+ public static String softDeleteSemanticModelVersionsBySchemaIds(
+ @Param("schemaIds") List<Long> schemaIds) {
+ return getProvider().softDeleteSemanticModelVersionsBySchemaIds(schemaIds);
+ }
+
+ /** Provides SQL for soft-deleting snapshots by catalog ID. */
+ public static String softDeleteSemanticModelVersionsByCatalogId(
+ @Param("catalogId") Long catalogId) {
+ return getProvider().softDeleteSemanticModelVersionsByCatalogId(catalogId);
+ }
+
+ /** Provides SQL for soft-deleting snapshots by metalake ID. */
+ public static String softDeleteSemanticModelVersionsByMetalakeId(
+ @Param("metalakeId") Long metalakeId) {
+ return
getProvider().softDeleteSemanticModelVersionsByMetalakeId(metalakeId);
+ }
+
+ /** Provides SQL for permanently deleting old snapshots. */
+ public static String deleteSemanticModelVersionsByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
+ return
getProvider().deleteSemanticModelVersionsByLegacyTimeline(legacyTimeline,
limit);
+ }
+
+ /** Provides SQL for selecting models that exceed the version retention
count. */
+ public static String selectSemanticModelVersionsByRetentionCount(
+ @Param("versionRetentionCount") Long versionRetentionCount) {
+ return
getProvider().selectSemanticModelVersionsByRetentionCount(versionRetentionCount);
+ }
+
+ /** Provides SQL for soft-deleting snapshots through a retention line. */
+ public static String softDeleteSemanticModelVersionsByRetentionLine(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("versionRetentionLine") long versionRetentionLine,
+ @Param("limit") int limit) {
+ return getProvider()
+ .softDeleteSemanticModelVersionsByRetentionLine(
+ semanticModelId, versionRetentionLine, limit);
+ }
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
index 00a72eb5b0..b7560254df 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
@@ -26,6 +26,7 @@ import
org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper;
import org.apache.gravitino.storage.relational.mapper.FunctionMetaMapper;
import org.apache.gravitino.storage.relational.mapper.MetalakeMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ModelMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
@@ -211,6 +212,9 @@ public class SchemaMetaBaseSQLProvider {
+ " UNION ALL SELECT 1 FROM "
+ TopicMetaMapper.TABLE_NAME
+ " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ + " UNION ALL SELECT 1 FROM "
+ + SemanticModelMetaMapper.TABLE_NAME
+ + " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ " LIMIT 1";
}
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
index f32a439959..66b57f8e59 100644
---
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
@@ -21,13 +21,14 @@ 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 java.util.List;
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. */
+/** Provides MySQL-compatible SQL for Semantic Model identity metadata. */
public class SemanticModelMetaBaseSQLProvider {
private static final String CURRENT_SNAPSHOT_COLUMNS =
@@ -43,6 +44,28 @@ public class SemanticModelMetaBaseSQLProvider {
+ " smvi.audit_info as version_audit_info,"
+ " smvi.deleted_at as version_deleted_at";
+ /** Returns SQL for listing current Semantic Model snapshots by schema ID. */
+ public String listSemanticModelPOsBySchemaId(@Param("schemaId") Long
schemaId) {
+ 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.deleted_at = 0 AND smvi.deleted_at = 0";
+ }
+
+ /** Returns SQL for listing current Semantic Model snapshots by qualified
schema name. */
+ public String listSemanticModelPOsByFullQualifiedName(
+ @Param("metalakeName") String metalakeName,
+ @Param("catalogName") String catalogName,
+ @Param("schemaName") String schemaName) {
+ return qualifiedNameSelect(false);
+ }
+
/** 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) {
@@ -109,12 +132,93 @@ public class SemanticModelMetaBaseSQLProvider {
+ " WHERE semantic_model_id = #{semanticModelId} AND deleted_at = 0
FOR UPDATE";
}
+ /** 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)";
+ }
+
+ /** Returns SQL for soft-deleting a Semantic Model identity with an
optimistic version check. */
+ public String softDeleteSemanticModelMetasBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("currentVersion") Integer currentVersion) {
+ return softDeleteBy(
+ "semantic_model_id = #{semanticModelId} AND current_version =
#{currentVersion}");
+ }
+
+ /** Returns SQL for soft-deleting Semantic Model identities by metalake ID.
*/
+ public String softDeleteSemanticModelMetasByMetalakeId(@Param("metalakeId")
Long metalakeId) {
+ return softDeleteBy("metalake_id = #{metalakeId}");
+ }
+
+ /** Returns SQL for soft-deleting Semantic Model identities by catalog ID. */
+ public String softDeleteSemanticModelMetasByCatalogId(@Param("catalogId")
Long catalogId) {
+ return softDeleteBy("catalog_id = #{catalogId}");
+ }
+
+ /** Returns SQL for soft-deleting Semantic Model identities by schema IDs. */
+ public String softDeleteSemanticModelMetasBySchemaIds(@Param("schemaIds")
List<Long> schemaIds) {
+ return "<script>UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}</foreach>) AND deleted_at = 0</script>";
+ }
+
+ /** Returns SQL for permanently deleting old Semantic Model identity rows. */
+ public String deleteSemanticModelMetasByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
+ return "DELETE FROM "
+ + TABLE_NAME
+ + " WHERE deleted_at > 0 AND deleted_at < #{legacyTimeline} LIMIT
#{limit}";
+ }
+
/** 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 qualifiedNameSelect(true);
+ }
+
+ private String qualifiedNameSelect(boolean filterBySemanticModelName) {
+ String semanticModelNameCondition =
+ filterBySemanticModelName ? " AND smm.semantic_model_name =
#{semanticModelName}" : "";
return """
SELECT
mm.metalake_id,
@@ -149,8 +253,7 @@ public class SemanticModelMetaBaseSQLProvider {
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}
+ %s smm ON sm.schema_id = smm.schema_id%s
AND smm.deleted_at = 0
LEFT JOIN
%s smvi ON smm.semantic_model_id = smvi.semantic_model_id
@@ -165,44 +268,17 @@ public class SemanticModelMetaBaseSQLProvider {
CatalogMetaMapper.TABLE_NAME,
SchemaMetaMapper.TABLE_NAME,
TABLE_NAME,
+ semanticModelNameCondition,
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) {
+ private String softDeleteBy(String condition) {
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)";
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE "
+ + condition
+ + " AND 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
index 8873f1c77e..c83286691b 100644
---
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
@@ -18,11 +18,12 @@
*/
package org.apache.gravitino.storage.relational.mapper.provider.base;
+import java.util.List;
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. */
+/** Provides MySQL-compatible SQL for Semantic Model version snapshots. */
public class SemanticModelVersionInfoBaseSQLProvider {
/** Returns SQL for inserting a Semantic Model version snapshot. */
@@ -42,4 +43,86 @@ public class SemanticModelVersionInfoBaseSQLProvider {
+ " #{semanticModelVersionInfo.properties},
#{semanticModelVersionInfo.auditInfo},"
+ " #{semanticModelVersionInfo.deletedAt})";
}
+
+ /** Returns SQL for selecting a Semantic Model version snapshot. */
+ public String selectSemanticModelVersionInfoBySemanticModelIdAndVersion(
+ @Param("semanticModelId") Long semanticModelId, @Param("version")
Integer version) {
+ return "SELECT id as id, metalake_id as metalakeId, catalog_id as
catalogId,"
+ + " schema_id as schemaId, semantic_model_id as semanticModelId,
version as version,"
+ + " semantic_model_name as semanticModelName,"
+ + " semantic_model_comment as semanticModelComment,"
+ + " semantic_model_definition as semanticModelDefinition, properties
as properties,"
+ + " audit_info as auditInfo, deleted_at as deletedAt FROM "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " WHERE semantic_model_id = #{semanticModelId}"
+ + " AND version = #{version} AND deleted_at = 0";
+ }
+
+ /** Returns SQL for soft-deleting all snapshots for a Semantic Model ID. */
+ public String softDeleteSemanticModelVersionsBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId) {
+ return softDeleteBy("semantic_model_id = #{semanticModelId}");
+ }
+
+ /** Returns SQL for soft-deleting Semantic Model snapshots by schema IDs. */
+ public String softDeleteSemanticModelVersionsBySchemaIds(
+ @Param("schemaIds") List<Long> schemaIds) {
+ return "<script>UPDATE "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}</foreach>) AND deleted_at = 0</script>";
+ }
+
+ /** Returns SQL for soft-deleting Semantic Model snapshots by catalog ID. */
+ public String softDeleteSemanticModelVersionsByCatalogId(@Param("catalogId")
Long catalogId) {
+ return softDeleteBy("catalog_id = #{catalogId}");
+ }
+
+ /** Returns SQL for soft-deleting Semantic Model snapshots by metalake ID. */
+ public String
softDeleteSemanticModelVersionsByMetalakeId(@Param("metalakeId") Long
metalakeId) {
+ return softDeleteBy("metalake_id = #{metalakeId}");
+ }
+
+ /** Returns SQL for permanently deleting old Semantic Model snapshots. */
+ public String deleteSemanticModelVersionsByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
+ return "DELETE FROM "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " WHERE deleted_at > 0 AND deleted_at < #{legacyTimeline} LIMIT
#{limit}";
+ }
+
+ /** Returns SQL for finding Semantic Models that exceed a version retention
count. */
+ public String selectSemanticModelVersionsByRetentionCount(
+ @Param("versionRetentionCount") Long versionRetentionCount) {
+ return "SELECT semantic_model_id as semanticModelId, MAX(version) as
version FROM "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " WHERE version > #{versionRetentionCount} AND deleted_at = 0"
+ + " GROUP BY semantic_model_id";
+ }
+
+ /** Returns SQL for soft-deleting snapshots through a per-model retention
line. */
+ public String softDeleteSemanticModelVersionsByRetentionLine(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("versionRetentionLine") long versionRetentionLine,
+ @Param("limit") int limit) {
+ return "UPDATE "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE semantic_model_id = #{semanticModelId}"
+ + " AND version <= #{versionRetentionLine} AND deleted_at = 0 LIMIT
#{limit}";
+ }
+
+ private String softDeleteBy(String condition) {
+ return "UPDATE "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE "
+ + condition
+ + " AND deleted_at = 0";
+ }
}
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
index 1261cfbcce..83457e5704 100644
---
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
@@ -18,7 +18,59 @@
*/
package org.apache.gravitino.storage.relational.mapper.provider.postgresql;
+import static
org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper.TABLE_NAME;
+
+import java.util.List;
import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelMetaBaseSQLProvider;
+import org.apache.ibatis.annotations.Param;
+
+/** Provides PostgreSQL SQL for Semantic Model identity metadata. */
+public class SemanticModelMetaPostgreSQLProvider extends
SemanticModelMetaBaseSQLProvider {
+
+ @Override
+ public String softDeleteSemanticModelMetasBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("currentVersion") Integer currentVersion) {
+ return softDeleteBy(
+ "semantic_model_id = #{semanticModelId} AND current_version =
#{currentVersion}");
+ }
+
+ @Override
+ public String softDeleteSemanticModelMetasByMetalakeId(@Param("metalakeId")
Long metalakeId) {
+ return softDeleteBy("metalake_id = #{metalakeId}");
+ }
+
+ @Override
+ public String softDeleteSemanticModelMetasByCatalogId(@Param("catalogId")
Long catalogId) {
+ return softDeleteBy("catalog_id = #{catalogId}");
+ }
+
+ @Override
+ public String softDeleteSemanticModelMetasBySchemaIds(@Param("schemaIds")
List<Long> schemaIds) {
+ return "<script>UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}</foreach>) AND deleted_at = 0</script>";
+ }
+
+ @Override
+ public String deleteSemanticModelMetasByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
+ return "DELETE FROM "
+ + TABLE_NAME
+ + " WHERE semantic_model_id IN (SELECT semantic_model_id FROM "
+ + TABLE_NAME
+ + " WHERE deleted_at > 0 AND deleted_at < #{legacyTimeline} LIMIT
#{limit})";
+ }
-/** Provides PostgreSQL SQL for Semantic Model create and load operations. */
-public class SemanticModelMetaPostgreSQLProvider extends
SemanticModelMetaBaseSQLProvider {}
+ private String softDeleteBy(String condition) {
+ return "UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE "
+ + condition
+ + " AND deleted_at = 0";
+ }
+}
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
index 9f78550b93..52643579ef 100644
---
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
@@ -18,8 +18,73 @@
*/
package org.apache.gravitino.storage.relational.mapper.provider.postgresql;
+import java.util.List;
+import
org.apache.gravitino.storage.relational.mapper.SemanticModelVersionInfoMapper;
import
org.apache.gravitino.storage.relational.mapper.provider.base.SemanticModelVersionInfoBaseSQLProvider;
+import org.apache.ibatis.annotations.Param;
-/** Provides PostgreSQL SQL for creating Semantic Model version snapshots. */
+/** Provides PostgreSQL SQL for Semantic Model version snapshots. */
public class SemanticModelVersionInfoPostgreSQLProvider
- extends SemanticModelVersionInfoBaseSQLProvider {}
+ extends SemanticModelVersionInfoBaseSQLProvider {
+
+ @Override
+ public String softDeleteSemanticModelVersionsBySemanticModelId(
+ @Param("semanticModelId") Long semanticModelId) {
+ return softDeleteBy("semantic_model_id = #{semanticModelId}");
+ }
+
+ @Override
+ public String softDeleteSemanticModelVersionsBySchemaIds(
+ @Param("schemaIds") List<Long> schemaIds) {
+ return "<script>UPDATE "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}</foreach>) AND deleted_at = 0</script>";
+ }
+
+ @Override
+ public String softDeleteSemanticModelVersionsByCatalogId(@Param("catalogId")
Long catalogId) {
+ return softDeleteBy("catalog_id = #{catalogId}");
+ }
+
+ @Override
+ public String
softDeleteSemanticModelVersionsByMetalakeId(@Param("metalakeId") Long
metalakeId) {
+ return softDeleteBy("metalake_id = #{metalakeId}");
+ }
+
+ @Override
+ public String deleteSemanticModelVersionsByLegacyTimeline(
+ @Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
+ return "DELETE FROM "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " WHERE id IN (SELECT id FROM "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " WHERE deleted_at > 0 AND deleted_at < #{legacyTimeline} LIMIT
#{limit})";
+ }
+
+ @Override
+ public String softDeleteSemanticModelVersionsByRetentionLine(
+ @Param("semanticModelId") Long semanticModelId,
+ @Param("versionRetentionLine") long versionRetentionLine,
+ @Param("limit") int limit) {
+ return "UPDATE "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE id IN (SELECT id FROM "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " WHERE semantic_model_id = #{semanticModelId}"
+ + " AND version <= #{versionRetentionLine}"
+ + " AND deleted_at = 0 LIMIT #{limit})";
+ }
+
+ private String softDeleteBy(String condition) {
+ return "UPDATE "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE "
+ + condition
+ + " AND deleted_at = 0";
+ }
+}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
index c80fe23411..8ef3740548 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
@@ -48,6 +48,8 @@ import
org.apache.gravitino.storage.relational.mapper.ModelVersionMetaMapper;
import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
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;
@@ -352,7 +354,15 @@ public class CatalogMetaService {
() ->
SessionUtils.doWithoutCommit(
ViewVersionInfoMapper.class,
- mapper ->
mapper.softDeleteViewVersionsByCatalogId(catalogId)));
+ mapper ->
mapper.softDeleteViewVersionsByCatalogId(catalogId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
mapper.softDeleteSemanticModelMetasByCatalogId(catalogId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
mapper.softDeleteSemanticModelVersionsByCatalogId(catalogId)));
} else {
SessionUtils.doMultipleWithCommit(
() -> {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
index ee811c7d44..c4eb17fe3f 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
@@ -54,6 +54,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;
@@ -336,7 +338,15 @@ public class MetalakeMetaService {
() ->
SessionUtils.doWithoutCommit(
ViewVersionInfoMapper.class,
- mapper ->
mapper.softDeleteViewVersionsByMetalakeId(metalakeId)));
+ mapper ->
mapper.softDeleteViewVersionsByMetalakeId(metalakeId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
mapper.softDeleteSemanticModelMetasByMetalakeId(metalakeId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
mapper.softDeleteSemanticModelVersionsByMetalakeId(metalakeId)));
} else {
SessionUtils.doMultipleWithCommit(
() -> {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
index 671af53f1c..7e6f3990af 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
@@ -55,6 +55,8 @@ import
org.apache.gravitino.storage.relational.mapper.ModelVersionMetaMapper;
import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
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;
@@ -359,7 +361,15 @@ public class SchemaMetaService {
() ->
SessionUtils.doWithoutCommit(
ViewVersionInfoMapper.class,
- mapper ->
mapper.softDeleteViewVersionsBySchemaIds(schemaIds.get())));
+ mapper ->
mapper.softDeleteViewVersionsBySchemaIds(schemaIds.get())),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
mapper.softDeleteSemanticModelMetasBySchemaIds(schemaIds.get())),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
mapper.softDeleteSemanticModelVersionsBySchemaIds(schemaIds.get())));
} else {
SessionUtils.doMultipleWithCommit(
() -> {
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
index acc13fbc07..4baffc78ee 100644
---
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
@@ -19,30 +19,45 @@
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.buildSemanticModelPO;
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.List;
import java.util.Locale;
import java.util.Objects;
+import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Function;
+import java.util.stream.Collectors;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.HasIdentifier;
import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.Namespace;
import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.exceptions.OptimisticLockException;
import org.apache.gravitino.meta.SemanticModelEntity;
import org.apache.gravitino.metrics.Monitored;
+import
org.apache.gravitino.storage.relational.EntityChangeLogNameIdentifierCodec;
+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.po.SemanticModelVersionInfoPO;
+import org.apache.gravitino.storage.relational.po.cache.OperateType;
import org.apache.gravitino.storage.relational.utils.ExceptionUtils;
import org.apache.gravitino.storage.relational.utils.SessionUtils;
import org.apache.gravitino.utils.NameIdentifierUtil;
+import org.apache.gravitino.utils.NamespaceUtil;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
-/** Provides relational create and load operations for Semantic Model
metadata. */
+/** Provides relational persistence operations for Semantic Model metadata. */
public class SemanticModelMetaService {
+ private static final Logger LOG =
LoggerFactory.getLogger(SemanticModelMetaService.class);
private static final SemanticModelMetaService INSTANCE = new
SemanticModelMetaService();
private final BasePOStorageOps<SemanticModelPO, SemanticModelMetaMapper> ops;
@@ -68,17 +83,33 @@ public class SemanticModelMetaService {
metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
baseMetricName = "getSemanticModelIdBySchemaIdAndName")
public Long getSemanticModelIdBySchemaIdAndName(long schemaId, String
semanticModelName) {
- Long semanticModelId =
+ SemanticModelPO semanticModelPO =
SessionUtils.getWithoutCommit(
SemanticModelMetaMapper.class,
- mapper -> mapper.selectSemanticModelIdBySchemaIdAndName(schemaId,
semanticModelName));
- if (semanticModelId == null) {
+ mapper -> ops.getPO(mapper, schemaId, semanticModelName));
+ if (semanticModelPO == null) {
throw new NoSuchEntityException(
NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
Entity.EntityType.SEMANTIC_MODEL.name().toLowerCase(Locale.ROOT),
semanticModelName);
}
- return semanticModelId;
+ return semanticModelPO.getSemanticModelId();
+ }
+
+ /**
+ * Lists current Semantic Models under a namespace.
+ *
+ * @param namespace The Semantic Model namespace.
+ * @return The current Semantic Model entities.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "listSemanticModelsByNamespace")
+ public List<SemanticModelEntity> listSemanticModelsByNamespace(Namespace
namespace) {
+ NamespaceUtil.checkSemanticModel(namespace);
+ return listSemanticModelPOs(namespace).stream()
+ .map(po -> fromSemanticModelPO(po, namespace))
+ .collect(Collectors.toList());
}
/**
@@ -156,11 +187,246 @@ public class SemanticModelMetaService {
}
}
+ /**
+ * Atomically creates a complete new snapshot and advances the current
version pointer.
+ *
+ * @param identifier The current Semantic Model identifier.
+ * @param updater The entity updater.
+ * @param <E> The internal entity type accepted by the updater.
+ * @return The updated Semantic Model entity.
+ * @throws IOException If persistence fails.
+ * @throws NoSuchEntityException If the Semantic Model is deleted or renamed
during the update.
+ * @throws OptimisticLockException If the internal transaction loses a
concurrent update race.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "updateSemanticModel")
+ public <E extends Entity & HasIdentifier> SemanticModelEntity
updateSemanticModel(
+ NameIdentifier identifier, Function<E, E> updater) throws IOException {
+ SemanticModelPO oldSemanticModelPO =
getSemanticModelPOByIdentifier(identifier);
+ SemanticModelEntity oldSemanticModelEntity =
+ fromSemanticModelPO(oldSemanticModelPO, identifier.namespace());
+ SemanticModelEntity newEntity = (SemanticModelEntity) updater.apply((E)
oldSemanticModelEntity);
+ Preconditions.checkArgument(
+ Objects.equals(oldSemanticModelEntity.id(), newEntity.id()),
+ "The updated Semantic Model entity id: %s should be same with the
entity id before: %s",
+ newEntity.id(),
+ oldSemanticModelEntity.id());
+ Preconditions.checkArgument(
+ oldSemanticModelEntity.namespace().equals(newEntity.namespace()),
+ "Semantic Model namespace cannot change from %s to %s",
+ oldSemanticModelEntity.namespace(),
+ newEntity.namespace());
+
+ AtomicInteger updateResult = new AtomicInteger(-1);
+ try {
+ SemanticModelPO newSemanticModelPO =
updateSemanticModelPO(oldSemanticModelPO, newEntity);
+ String metalakeName = identifier.namespace().level(0);
+ String catalogName = identifier.namespace().level(1);
+ String schemaName = identifier.namespace().level(2);
+ String oldFullName =
+ EntityChangeLogNameIdentifierCodec.encode(
+ NameIdentifierUtil.ofSemanticModel(
+ metalakeName,
+ catalogName,
+ schemaName,
+ oldSemanticModelPO.getSemanticModelName()));
+ boolean isRenamed =
+ !Objects.equals(
+ oldSemanticModelPO.getSemanticModelName(),
newSemanticModelPO.getSemanticModelName());
+
+ SessionUtils.doMultipleWithCommit(
+ // The Semantic Model and its parent were read before this
transaction started. Lock the
+ // observed parent again before either identity or version writes,
so a schema cascade
+ // cannot finish its cleanup and then let this update recreate child
state below it.
+ () ->
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ identifier,
+ oldSemanticModelPO.getSchemaId(),
+ oldSemanticModelPO.getCatalogId(),
+ oldSemanticModelPO.getMetalakeId()),
+ () -> {
+ updateResult.set(
+ SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper -> ops.updatePO(mapper, newSemanticModelPO,
oldSemanticModelPO)));
+ if (updateResult.get() == 0) {
+ throw semanticModelWriteFailure(identifier, oldSemanticModelPO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
+ mapper.insertSemanticModelVersionInfo(
+ newSemanticModelPO.getSemanticModelVersionInfoPO())),
+ () -> {
+ if (isRenamed && updateResult.get() > 0) {
+ SessionUtils.doWithoutCommit(
+ EntityChangeLogMapper.class,
+ mapper ->
+ mapper.insertEntityChange(
+ metalakeName,
+ Entity.EntityType.SEMANTIC_MODEL.name(),
+ oldFullName,
+ OperateType.ALTER));
+ }
+ });
+ return newEntity;
+ } catch (OptimisticLockException optimisticLockException) {
+ throw optimisticLockException;
+ } catch (RuntimeException re) {
+ ExceptionUtils.checkSQLException(
+ re, Entity.EntityType.SEMANTIC_MODEL,
newEntity.nameIdentifier().toString());
+ throw re;
+ }
+ }
+
+ /**
+ * Soft-deletes a Semantic Model identity and all of its version snapshots.
+ *
+ * @param identifier The Semantic Model identifier.
+ * @return {@code true} when an active identity was deleted.
+ * @throws NoSuchEntityException If the Semantic Model is already deleted or
renamed.
+ * @throws OptimisticLockException If the internal transaction loses a
concurrent update race.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "deleteSemanticModel")
+ public boolean deleteSemanticModel(NameIdentifier identifier) {
+ SemanticModelPO semanticModelPO =
getSemanticModelPOByIdentifier(identifier);
+ deleteSemanticModelWithVersion(identifier, semanticModelPO);
+ return true;
+ }
+
+ /**
+ * Permanently deletes soft-deleted Semantic Model identities and snapshots
older than a timeline.
+ *
+ * @param legacyTimeline The exclusive deletion timeline in epoch
milliseconds.
+ * @param limit The maximum rows to delete from each table.
+ * @return The total number of deleted rows.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "deleteSemanticModelMetasByLegacyTimeline")
+ public int deleteSemanticModelMetasByLegacyTimeline(Long legacyTimeline, int
limit) {
+ int versionDeletedCount =
+ SessionUtils.doWithCommitAndFetchResult(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
mapper.deleteSemanticModelVersionsByLegacyTimeline(legacyTimeline, limit));
+ int metaDeletedCount =
+ SessionUtils.doWithCommitAndFetchResult(
+ SemanticModelMetaMapper.class,
+ mapper ->
mapper.deleteSemanticModelMetasByLegacyTimeline(legacyTimeline, limit));
+ return versionDeletedCount + metaDeletedCount;
+ }
+
+ /**
+ * Soft-deletes old Semantic Model snapshots beyond the configured retention
count.
+ *
+ * @param versionRetentionCount The number of latest versions to retain per
Semantic Model.
+ * @param limit The maximum versions to delete for each Semantic Model.
+ * @return The number of snapshots soft-deleted.
+ */
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "deleteSemanticModelVersionsByRetentionCount")
+ public int deleteSemanticModelVersionsByRetentionCount(Long
versionRetentionCount, int limit) {
+ List<SemanticModelVersionInfoPO> currentVersions =
+ SessionUtils.getWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
mapper.selectSemanticModelVersionsByRetentionCount(versionRetentionCount));
+
+ int totalDeletedCount = 0;
+ for (SemanticModelVersionInfoPO currentVersion : currentVersions) {
+ long versionRetentionLine = currentVersion.version() -
versionRetentionCount;
+ int deletedCount =
+ SessionUtils.doWithCommitAndFetchResult(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
+ mapper.softDeleteSemanticModelVersionsByRetentionLine(
+ currentVersion.semanticModelId(), versionRetentionLine,
limit));
+ totalDeletedCount += deletedCount;
+ LOG.info(
+ "Soft delete Semantic Model versions count: {} through retention
line: {},"
+ + " current Semantic Model id and version: <{}, {}>.",
+ deletedCount,
+ versionRetentionLine,
+ currentVersion.semanticModelId(),
+ currentVersion.version());
+ }
+ return totalDeletedCount;
+ }
+
/** Returns the persistent-object operations used by this service. */
public BasePOStorageOps<SemanticModelPO, SemanticModelMetaMapper> ops() {
return ops;
}
+ /**
+ * Deletes the observed Semantic Model and its snapshots in one transaction.
+ *
+ * <p>Package-private access lets concurrency tests submit a stale snapshot
through the same
+ * compare-and-set path as public deletion.
+ */
+ void deleteSemanticModelWithVersion(
+ NameIdentifier identifier, SemanticModelPO observedSemanticModelPO) {
+ Long semanticModelId = observedSemanticModelPO.getSemanticModelId();
+ String metalakeName = identifier.namespace().level(0);
+ String semanticModelFullName =
+ EntityChangeLogNameIdentifierCodec.encode(
+ NameIdentifierUtil.ofSemanticModel(
+ metalakeName,
+ identifier.namespace().level(1),
+ identifier.namespace().level(2),
+ observedSemanticModelPO.getSemanticModelName()));
+ // TODO: Soft-delete Semantic Model owner, tag, and securable-object
relations in this
+ // transaction and in the metalake/catalog/schema cascade delete paths.
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ OccWriteSupport.deleteWithVersion(
+ () ->
+ SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
+
mapper.softDeleteSemanticModelMetasBySemanticModelId(
+ semanticModelId,
observedSemanticModelPO.getCurrentVersion())),
+ () -> semanticModelWriteFailure(identifier,
observedSemanticModelPO)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
mapper.softDeleteSemanticModelVersionsBySemanticModelId(semanticModelId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ EntityChangeLogMapper.class,
+ mapper ->
+ mapper.insertEntityChange(
+ metalakeName,
+ Entity.EntityType.SEMANTIC_MODEL.name(),
+ semanticModelFullName,
+ OperateType.DROP)));
+ }
+
+ private SemanticModelPO updateSemanticModelPO(
+ SemanticModelPO oldSemanticModelPO, SemanticModelEntity newEntity) {
+ int previousVersion =
+ Math.max(oldSemanticModelPO.getCurrentVersion(),
oldSemanticModelPO.getLastVersion());
+ Preconditions.checkState(
+ previousVersion < Integer.MAX_VALUE,
+ "Semantic Model %s has exhausted the version range",
+ oldSemanticModelPO.getSemanticModelId());
+ int newVersion = previousVersion + 1;
+ SemanticModelPO.SemanticModelPOBuilder builder =
+ SemanticModelPO.builder()
+ .withMetalakeId(oldSemanticModelPO.getMetalakeId())
+ .withCatalogId(oldSemanticModelPO.getCatalogId())
+ .withSchemaId(oldSemanticModelPO.getSchemaId())
+ .withCurrentVersion(newVersion)
+ .withLastVersion(newVersion);
+ return buildSemanticModelPO(newEntity, builder, newVersion);
+ }
+
private static SemanticModelPO semanticModelForOverwrite(
SemanticModelPO source, SemanticModelPO persistedPO) {
int previousVersion = Math.max(persistedPO.getCurrentVersion(),
persistedPO.getLastVersion());
@@ -202,7 +468,7 @@ public class SemanticModelMetaService {
.build();
}
- private SemanticModelPO getSemanticModelPOByIdentifier(NameIdentifier
identifier) {
+ SemanticModelPO getSemanticModelPOByIdentifier(NameIdentifier identifier) {
NameIdentifierUtil.checkSemanticModel(identifier);
SemanticModelPO semanticModelPO =
SessionUtils.getWithoutCommit(
@@ -239,4 +505,11 @@ public class SemanticModelMetaService {
&& Objects.equals(
current.getMetalakeId(),
observedSemanticModelPO.getMetalakeId()));
}
+
+ private List<SemanticModelPO> listSemanticModelPOs(Namespace namespace) {
+ return SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
+ POStorageReadRouting.listPOs(mapper, namespace, ops,
Entity.EntityType.SEMANTIC_MODEL));
+ }
}
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
index a2877328f7..e1e2ac5f31 100644
---
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
@@ -19,7 +19,9 @@
package org.apache.gravitino.storage.relational.service;
import com.google.common.base.Preconditions;
+import java.util.List;
import java.util.Locale;
+import java.util.stream.Collectors;
import org.apache.gravitino.Entity;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
@@ -27,7 +29,7 @@ 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. */
+/** Provides relational persistent-object operations for Semantic Models. */
public class SemanticModelPOStorageOps
extends BasePOStorageOps<SemanticModelPO, SemanticModelMetaMapper> {
@@ -49,6 +51,12 @@ public class SemanticModelPOStorageOps
mapper.insertSemanticModelMeta(semanticModelPO);
}
+ @Override
+ public Integer updatePO(
+ SemanticModelMetaMapper mapper, SemanticModelPO newPO, SemanticModelPO
oldPO) {
+ return mapper.updateSemanticModelMeta(newPO, oldPO);
+ }
+
@Override
public SemanticModelPO getPO(
SemanticModelMetaMapper mapper, Long parentId, String semanticModelName)
{
@@ -80,6 +88,32 @@ public class SemanticModelPOStorageOps
return po;
}
+ @Override
+ public List<SemanticModelPO> listPOs(SemanticModelMetaMapper mapper, Long
parentId) {
+ return mapper.listSemanticModelPOsBySchemaId(parentId);
+ }
+
+ @Override
+ public List<SemanticModelPO> listPOsByNSFullName(
+ SemanticModelMetaMapper mapper, Namespace namespace) {
+ List<SemanticModelPO> pos =
+ mapper.listSemanticModelPOsByFullQualifiedName(
+ namespace.level(0), namespace.level(1), namespace.level(2));
+ if (pos.isEmpty()) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.CATALOG.name().toLowerCase(Locale.ROOT),
+ namespace.level(1));
+ }
+ if (pos.get(0).getSchemaId() == null) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.SCHEMA.name().toLowerCase(Locale.ROOT),
+ namespace.level(2));
+ }
+ return pos.stream().filter(po -> po.getSemanticModelId() !=
null).collect(Collectors.toList());
+ }
+
@Override
public boolean supportsParentIdRelationalRead() {
return true;
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
b/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
index 1e4fe8c884..fa1c2fc554 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
@@ -252,6 +252,10 @@ public abstract class TestJDBCBackend {
tableName = "function_meta";
idColumnName = "function_id";
break;
+ case SEMANTIC_MODEL:
+ tableName = "semantic_model_meta";
+ idColumnName = "semantic_model_id";
+ break;
default:
throw new IllegalArgumentException("Unsupported entity type: " +
entityType);
}
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
index 566c51fea3..c297a4efc0 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/TestSemanticModelJDBCBackend.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/TestSemanticModelJDBCBackend.java
@@ -19,6 +19,7 @@
package org.apache.gravitino.storage.relational;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -32,6 +33,7 @@ import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
+import java.time.Instant;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
@@ -47,6 +49,7 @@ 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.exceptions.NonEmptyEntityException;
import org.apache.gravitino.meta.NamespacedEntityId;
import org.apache.gravitino.meta.SemanticModelEntity;
import org.apache.gravitino.semantic.AIContext;
@@ -74,7 +77,7 @@ 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}. */
+/** Tests Semantic Model persistence and parent lifecycle behavior through
{@link JDBCBackend}. */
public class TestSemanticModelJDBCBackend extends TestJDBCBackend {
@TestTemplate
@@ -300,6 +303,8 @@ public class TestSemanticModelJDBCBackend extends
TestJDBCBackend {
SemanticModelPO byFullName =
readSemanticModelPO(semanticModel.nameIdentifier(), false);
assertEquals(semanticModel.id(), byParentId.getSemanticModelId());
assertEquals(semanticModel.id(), byFullName.getSemanticModelId());
+ assertEquals(List.of(persisted), listSemanticModelPOs(namespace, true));
+ assertEquals(List.of(persisted), listSemanticModelPOs(namespace, false));
RelationalEntityStoreIdResolver resolver = new
RelationalEntityStoreIdResolver();
NamespacedEntityId resolved =
@@ -309,17 +314,38 @@ public class TestSemanticModelJDBCBackend extends
TestJDBCBackend {
String metalake = namespace.level(0);
String catalog = namespace.level(1);
String schema = namespace.level(2);
+ String emptySchema = "empty_schema";
+ createAndInsertSchema(metalake, catalog, emptySchema);
+ assertTrue(
+ listSemanticModelPOs(NamespaceUtil.ofSemanticModel(metalake, catalog,
emptySchema), false)
+ .isEmpty());
+
+ Namespace missingCatalogNamespace =
+ NamespaceUtil.ofSemanticModel(metalake, "missing_catalog", schema);
+ NoSuchEntityException missingCatalog =
+ assertThrows(
+ NoSuchEntityException.class,
+ () -> listSemanticModelPOs(missingCatalogNamespace, false));
+ assertEquals(
+ String.format(NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE, "catalog",
"missing_catalog"),
+ missingCatalog.getMessage());
+
+ Namespace missingSchemaNamespace =
+ NamespaceUtil.ofSemanticModel(metalake, catalog, "missing_schema");
+ NoSuchEntityException missingSchema =
+ assertThrows(
+ NoSuchEntityException.class, () ->
listSemanticModelPOs(missingSchemaNamespace, false));
+ assertEquals(
+ String.format(NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE, "schema",
"missing_schema"),
+ missingSchema.getMessage());
+
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()));
+ NameIdentifier.of(missingCatalogNamespace, semanticModel.name()),
+ NameIdentifier.of(missingSchemaNamespace, semanticModel.name()));
for (NameIdentifier missing : missingParents) {
assertThrows(NoSuchEntityException.class, () ->
readSemanticModelPO(missing, true));
assertThrows(NoSuchEntityException.class, () ->
readSemanticModelPO(missing, false));
@@ -390,6 +416,111 @@ public class TestSemanticModelJDBCBackend extends
TestJDBCBackend {
assertEquals(0, countEntityChanges());
}
+ @TestTemplate
+ public void testListUpdateDeleteAndGarbageCollection() throws IOException {
+ Namespace namespace = createParents("lifecycle");
+ SemanticModelEntity original =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "sales_model",
+ false,
+ ImmutableMap.of("domain", "sales"));
+ backend.insert(original, false);
+
+ assertEquals(
+ List.of(original), backend.list(namespace,
Entity.EntityType.SEMANTIC_MODEL, true));
+ assertEquals(
+ List.of(original),
+ backend.batchGet(
+ List.of(original.nameIdentifier(), NameIdentifier.of(namespace,
"missing_model")),
+ Entity.EntityType.SEMANTIC_MODEL));
+
+ SemanticModelEntity renamed =
+ semanticModel(
+ original.id(),
+ namespace,
+ "renamed_sales_model",
+ false,
+ ImmutableMap.of("domain", "sales_v2"));
+ assertEquals(
+ renamed,
+ backend.update(
+ original.nameIdentifier(), Entity.EntityType.SEMANTIC_MODEL,
ignored -> renamed));
+ assertFalse(backend.exists(original.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertEquals(renamed, backend.get(renamed.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+
+ assertTrue(backend.delete(renamed.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL, false));
+ assertFalse(backend.exists(renamed.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertTrue(legacyRecordExistsInDB(renamed.id(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertTrue(allVersionRowsAreSoftDeleted(renamed.id()));
+
+ backend.hardDeleteLegacyData(
+ Entity.EntityType.SEMANTIC_MODEL, Instant.now().toEpochMilli() + 1000);
+ assertFalse(legacyRecordExistsInDB(renamed.id(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertEquals(0, countRows("semantic_model_version_info", renamed.id()));
+ }
+
+ @TestTemplate
+ public void testSchemaNonCascadeAndCascadeLifecycle() throws IOException {
+ Namespace namespace = createParents("schema_cascade");
+ SemanticModelEntity semanticModel =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "schema_child_model",
+ false,
+ ImmutableMap.of());
+ backend.insert(semanticModel, false);
+
+ NameIdentifier schemaIdentifier =
+ NameIdentifier.of(namespace.level(0), namespace.level(1),
namespace.level(2));
+ assertThrows(
+ NonEmptyEntityException.class,
+ () -> backend.delete(schemaIdentifier, Entity.EntityType.SCHEMA,
false));
+ assertTrue(backend.exists(semanticModel.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+
+ assertTrue(backend.delete(schemaIdentifier, Entity.EntityType.SCHEMA,
true));
+ assertFalse(backend.exists(semanticModel.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertTrue(legacyRecordExistsInDB(semanticModel.id(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertTrue(allVersionRowsAreSoftDeleted(semanticModel.id()));
+ }
+
+ @TestTemplate
+ public void testCatalogAndMetalakeCascadeLifecycle() throws IOException {
+ Namespace catalogNamespace = createParents("catalog_cascade");
+ SemanticModelEntity catalogChild =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ catalogNamespace,
+ "catalog_child_model",
+ false,
+ ImmutableMap.of());
+ backend.insert(catalogChild, false);
+
+ NameIdentifier catalogIdentifier =
+ NameIdentifier.of(catalogNamespace.level(0),
catalogNamespace.level(1));
+ assertTrue(backend.delete(catalogIdentifier, Entity.EntityType.CATALOG,
true));
+ assertFalse(backend.exists(catalogChild.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertTrue(allVersionRowsAreSoftDeleted(catalogChild.id()));
+
+ Namespace metalakeNamespace = createParents("metalake_cascade");
+ SemanticModelEntity metalakeChild =
+ semanticModel(
+ RandomIdGenerator.INSTANCE.nextId(),
+ metalakeNamespace,
+ "metalake_child_model",
+ false,
+ ImmutableMap.of());
+ backend.insert(metalakeChild, false);
+
+ assertTrue(
+ backend.delete(
+ NameIdentifier.of(metalakeNamespace.level(0)),
Entity.EntityType.METALAKE, true));
+ assertFalse(backend.exists(metalakeChild.nameIdentifier(),
Entity.EntityType.SEMANTIC_MODEL));
+ assertTrue(allVersionRowsAreSoftDeleted(metalakeChild.id()));
+ }
+
private Namespace createParents(String prefix) throws IOException {
String suffix = Long.toUnsignedString(RandomIdGenerator.INSTANCE.nextId());
String metalakeName = prefix + "_metalake_" + suffix;
@@ -465,6 +596,23 @@ public class TestSemanticModelJDBCBackend extends
TestJDBCBackend {
throw new IllegalArgumentException("Unsupported Semantic Model table: "
+ tableName);
}
String sql = String.format("SELECT count(*) FROM %s WHERE
semantic_model_id = ?", tableName);
+ return queryCount(sql, semanticModelId);
+ }
+
+ private boolean allVersionRowsAreSoftDeleted(Long semanticModelId) {
+ int total = countRows(SemanticModelVersionInfoMapper.TABLE_NAME,
semanticModelId);
+ if (total == 0) {
+ return false;
+ }
+ return total
+ == queryCount(
+ "SELECT count(*) FROM "
+ + SemanticModelVersionInfoMapper.TABLE_NAME
+ + " WHERE semantic_model_id = ? AND deleted_at > 0",
+ semanticModelId);
+ }
+
+ private int queryCount(String sql, Long semanticModelId) {
try (SqlSession sqlSession =
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
Connection connection = sqlSession.getConnection();
@@ -500,6 +648,18 @@ public class TestSemanticModelJDBCBackend extends
TestJDBCBackend {
cacheEnabled));
}
+ private List<SemanticModelPO> listSemanticModelPOs(Namespace namespace,
boolean cacheEnabled) {
+ return SessionUtils.getWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
+ POStorageReadRouting.listPOs(
+ mapper,
+ namespace,
+ SemanticModelMetaService.getInstance().ops(),
+ Entity.EntityType.SEMANTIC_MODEL,
+ cacheEnabled));
+ }
+
private List<Integer> activeSnapshotVersions(Long semanticModelId) {
String sql =
String.format(
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestSchemaMetaBaseSQLProvider.java
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestSchemaMetaBaseSQLProvider.java
index 6cedd927d8..1cba5d3b57 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestSchemaMetaBaseSQLProvider.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestSchemaMetaBaseSQLProvider.java
@@ -23,6 +23,7 @@ import java.util.List;
import org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper;
import org.apache.gravitino.storage.relational.mapper.FunctionMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ModelMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
@@ -43,7 +44,8 @@ class TestSchemaMetaBaseSQLProvider {
FilesetMetaMapper.META_TABLE_NAME,
FunctionMetaMapper.TABLE_NAME,
ModelMetaMapper.TABLE_NAME,
- TopicMetaMapper.TABLE_NAME);
+ TopicMetaMapper.TABLE_NAME,
+ SemanticModelMetaMapper.TABLE_NAME);
childTables.forEach(
tableName ->
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSemanticModelMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSemanticModelMetaService.java
new file mode 100644
index 0000000000..5bc2f6888c
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSemanticModelMetaService.java
@@ -0,0 +1,687 @@
+/*
+ * 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.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+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.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.time.Instant;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+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 java.util.concurrent.TimeoutException;
+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.exceptions.OptimisticLockException;
+import org.apache.gravitino.integration.test.util.GravitinoITUtils;
+import org.apache.gravitino.meta.SchemaEntity;
+import org.apache.gravitino.meta.SemanticModelEntity;
+import org.apache.gravitino.semantic.AIContext;
+import org.apache.gravitino.semantic.CustomExtension;
+import org.apache.gravitino.semantic.Dataset;
+import org.apache.gravitino.semantic.SemanticModelDefinition;
+import org.apache.gravitino.storage.RandomIdGenerator;
+import
org.apache.gravitino.storage.relational.EntityChangeLogNameIdentifierCodec;
+import org.apache.gravitino.storage.relational.TestJDBCBackend;
+import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
+import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.SemanticModelMetaMapper;
+import
org.apache.gravitino.storage.relational.mapper.SemanticModelVersionInfoMapper;
+import org.apache.gravitino.storage.relational.po.SchemaPO;
+import org.apache.gravitino.storage.relational.po.SemanticModelPO;
+import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
+import org.apache.gravitino.storage.relational.po.cache.OperateType;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.gravitino.storage.relational.utils.SessionUtils;
+import org.apache.gravitino.utils.NameIdentifierUtil;
+import org.apache.gravitino.utils.NamespaceUtil;
+import org.apache.ibatis.session.SqlSession;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.TestTemplate;
+
+public class TestSemanticModelMetaService extends TestJDBCBackend {
+
+ private final String metalakeName =
GravitinoITUtils.genRandomName("tst_semantic_metalake");
+ private final String catalogName =
GravitinoITUtils.genRandomName("tst_semantic_catalog");
+ private final String schemaName =
GravitinoITUtils.genRandomName("tst_semantic_schema");
+
+ private Namespace namespace;
+ private long schemaId;
+
+ @BeforeEach
+ public void prepare() throws IOException {
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+ SchemaEntity schema = createAndInsertSchema(metalakeName, catalogName,
schemaName);
+ schemaId = schema.id();
+ namespace = NamespaceUtil.ofSemanticModel(metalakeName, catalogName,
schemaName);
+ }
+
+ @TestTemplate
+ public void testInsertLoadListOverwriteAndIdLookup() throws IOException {
+ String semanticModelName = GravitinoITUtils.genRandomName("sales_model");
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ semanticModelName,
+ "orders",
+ "initial comment",
+ "initial");
+
+ SemanticModelMetaService.getInstance().insertSemanticModel(original,
false);
+
+ assertEquals(
+ original.id(),
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelIdBySchemaIdAndName(schemaId, semanticModelName));
+ assertEquals(
+ original,
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(original.nameIdentifier()));
+ List<SemanticModelEntity> listed =
+
SemanticModelMetaService.getInstance().listSemanticModelsByNamespace(namespace);
+ assertEquals(1, listed.size());
+ assertEquals(original, listed.get(0));
+
+ SemanticModelEntity duplicate =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ semanticModelName,
+ "duplicate_orders",
+ "duplicate comment",
+ "duplicate");
+ assertThrows(
+ EntityAlreadyExistsException.class,
+ () ->
SemanticModelMetaService.getInstance().insertSemanticModel(duplicate, false));
+
+ SemanticModelEntity overwritten =
+ semanticModelEntity(
+ original.id(),
+ semanticModelName,
+ "orders_overwritten",
+ "overwritten comment",
+ "overwritten");
+ SemanticModelMetaService.getInstance().insertSemanticModel(overwritten,
true);
+ assertEquals(
+ overwritten,
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(overwritten.nameIdentifier()));
+ assertEquals(2, listSemanticModelVersions(original.id()).size());
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelIdBySchemaIdAndName(schemaId,
"missing_model"));
+ }
+
+ @TestTemplate
+ public void testUpdateCreatesFullSnapshotAndRenameChangeLog() throws
IOException {
+ String oldName = GravitinoITUtils.genRandomName("sales_model_old");
+ String newName = GravitinoITUtils.genRandomName("sales_model_new");
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(), oldName, "orders_v1", "v1
comment", "v1");
+ SemanticModelMetaService.getInstance().insertSemanticModel(original,
false);
+ long lastChangeId = maxEntityChangeId();
+
+ SemanticModelEntity expected =
+ semanticModelEntity(original.id(), newName, "orders_v2", "v2 comment",
"v2");
+ SemanticModelEntity updated =
+ SemanticModelMetaService.getInstance()
+ .updateSemanticModel(original.nameIdentifier(), ignored ->
expected);
+
+ assertEquals(expected, updated);
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(original.nameIdentifier()));
+ assertEquals(
+ expected,
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(expected.nameIdentifier()));
+ assertEquals(
+ original.id(),
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelIdBySchemaIdAndName(schemaId, newName));
+
+ Map<Integer, VersionState> versions =
listSemanticModelVersions(original.id());
+ assertEquals(2, versions.size());
+ assertEquals(oldName, versions.get(1).name);
+ assertEquals("v1 comment", versions.get(1).comment);
+ assertTrue(versions.get(1).definition.contains("orders_v1"));
+ assertEquals(newName, versions.get(2).name);
+ assertEquals("v2 comment", versions.get(2).comment);
+ assertTrue(versions.get(2).definition.contains("orders_v2"));
+ assertEquals(0L, versions.get(1).deletedAt);
+ assertEquals(0L, versions.get(2).deletedAt);
+
+ long[] identityVersions = semanticModelIdentityVersions(original.id());
+ assertEquals(2L, identityVersions[0]);
+ assertEquals(2L, identityVersions[1]);
+ assertEntityChange(lastChangeId, oldName, OperateType.ALTER);
+ }
+
+ @TestTemplate
+ public void testDottedNameRenameAndDropChangeLog() throws IOException {
+ String oldName = GravitinoITUtils.genRandomName("sales_model") + ".old";
+ String newName = GravitinoITUtils.genRandomName("sales_model") + ".new";
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(), oldName, "orders",
"original", "v1");
+ SemanticModelMetaService.getInstance().insertSemanticModel(original,
false);
+ long lastChangeId = maxEntityChangeId();
+
+ SemanticModelEntity renamed =
+ semanticModelEntity(original.id(), newName, "orders", "renamed", "v2");
+ SemanticModelMetaService.getInstance()
+ .updateSemanticModel(original.nameIdentifier(), ignored -> renamed);
+ assertEntityChange(lastChangeId, oldName, OperateType.ALTER);
+
+ lastChangeId = maxEntityChangeId();
+ assertTrue(
+
SemanticModelMetaService.getInstance().deleteSemanticModel(renamed.nameIdentifier()));
+ assertEntityChange(lastChangeId, newName, OperateType.DROP);
+ }
+
+ @TestTemplate
+ public void testUpdateRejectsNamespaceChange() throws IOException {
+ String semanticModelName =
GravitinoITUtils.genRandomName("namespace_change_model");
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ semanticModelName,
+ "orders",
+ "original comment",
+ "original");
+ SemanticModelMetaService.getInstance().insertSemanticModel(original,
false);
+
+ String targetSchemaName =
GravitinoITUtils.genRandomName("target_semantic_schema");
+ createAndInsertSchema(metalakeName, catalogName, targetSchemaName);
+ Namespace targetNamespace =
+ NamespaceUtil.ofSemanticModel(metalakeName, catalogName,
targetSchemaName);
+ SemanticModelEntity moved =
+ SemanticModelEntity.builder()
+ .withId(original.id())
+ .withName(original.name())
+ .withNamespace(targetNamespace)
+ .withComment(original.comment())
+ .withDefinition(original.definition())
+ .withProperties(original.properties())
+ .withAuditInfo(original.auditInfo())
+ .build();
+
+ IllegalArgumentException exception =
+ assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ SemanticModelMetaService.getInstance()
+ .updateSemanticModel(original.nameIdentifier(), ignored ->
moved));
+ assertEquals(
+ "Semantic Model namespace cannot change from " + namespace + " to " +
targetNamespace,
+ exception.getMessage());
+ assertEquals(
+ original,
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(original.nameIdentifier()));
+ assertTrue(
+ SemanticModelMetaService.getInstance()
+ .listSemanticModelsByNamespace(targetNamespace)
+ .isEmpty());
+ assertEquals(1L, semanticModelIdentityVersions(original.id())[0]);
+ assertEquals(1, listSemanticModelVersions(original.id()).size());
+ }
+
+ @TestTemplate
+ public void testInsertWaitsForConcurrentSchemaCascadeDelete() throws
Exception {
+ String semanticModelName =
GravitinoITUtils.genRandomName("concurrent_create_model");
+ SemanticModelEntity semanticModel =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(), semanticModelName, "orders",
"comment", "create");
+ SchemaPO observedSchemaPO = getObservedSchemaPO();
+ CountDownLatch schemaDeleteLocked = new CountDownLatch(1);
+ CountDownLatch allowDeleteCommit = new CountDownLatch(1);
+ CountDownLatch insertStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> deleteResult =
+ startSchemaCascadeDelete(observedSchemaPO, schemaDeleteLocked,
allowDeleteCommit, executor);
+
+ try {
+ assertTrue(schemaDeleteLocked.await(30, TimeUnit.SECONDS));
+ Future<Throwable> insertResult =
+ executor.submit(
+ () -> {
+ insertStarted.countDown();
+ try {
+
SemanticModelMetaService.getInstance().insertSemanticModel(semanticModel,
false);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ assertTrue(insertStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> insertResult.get(500,
TimeUnit.MILLISECONDS));
+
+ allowDeleteCommit.countDown();
+ assertNull(deleteResult.get(30, TimeUnit.SECONDS));
+ assertInstanceOf(NoSuchEntityException.class, insertResult.get(30,
TimeUnit.SECONDS));
+ } finally {
+ allowDeleteCommit.countDown();
+ executor.shutdownNow();
+ }
+
+ assertEquals(0, countSemanticModelIdentities(semanticModel.id()));
+ assertTrue(listSemanticModelVersions(semanticModel.id()).isEmpty());
+ }
+
+ @TestTemplate
+ public void testUpdateWaitsForConcurrentSchemaCascadeDelete() throws
Exception {
+ String semanticModelName =
GravitinoITUtils.genRandomName("concurrent_update_model");
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ semanticModelName,
+ "orders_v1",
+ "v1 comment",
+ "v1");
+ SemanticModelMetaService.getInstance().insertSemanticModel(original,
false);
+ SemanticModelEntity updated =
+ semanticModelEntity(original.id(), semanticModelName, "orders_v2", "v2
comment", "v2");
+ SchemaPO observedSchemaPO = getObservedSchemaPO();
+ CountDownLatch schemaDeleteLocked = new CountDownLatch(1);
+ CountDownLatch allowDeleteCommit = new CountDownLatch(1);
+ CountDownLatch updateStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> deleteResult =
+ startSchemaCascadeDelete(observedSchemaPO, schemaDeleteLocked,
allowDeleteCommit, executor);
+
+ try {
+ assertTrue(schemaDeleteLocked.await(30, TimeUnit.SECONDS));
+ Future<Throwable> updateResult =
+ executor.submit(
+ () -> {
+ updateStarted.countDown();
+ try {
+ SemanticModelMetaService.getInstance()
+ .updateSemanticModel(original.nameIdentifier(), ignored
-> updated);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ assertTrue(updateStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> updateResult.get(500,
TimeUnit.MILLISECONDS));
+
+ allowDeleteCommit.countDown();
+ assertNull(deleteResult.get(30, TimeUnit.SECONDS));
+ assertInstanceOf(NoSuchEntityException.class, updateResult.get(30,
TimeUnit.SECONDS));
+ } finally {
+ allowDeleteCommit.countDown();
+ executor.shutdownNow();
+ }
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(original.nameIdentifier()));
+ Map<Integer, VersionState> versions =
listSemanticModelVersions(original.id());
+ assertEquals(1, versions.size());
+ assertTrue(versions.get(1).deletedAt > 0L);
+ }
+
+ @TestTemplate
+ public void testOptimisticUpdateRollsBackInsertedSnapshot() throws
IOException {
+ String semanticModelName =
GravitinoITUtils.genRandomName("optimistic_model");
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ semanticModelName,
+ "orders_v1",
+ "v1 comment",
+ "v1");
+ SemanticModelMetaService.getInstance().insertSemanticModel(original,
false);
+ SemanticModelEntity expected =
+ semanticModelEntity(original.id(), semanticModelName, "orders_v2", "v2
comment", "v2");
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ SemanticModelMetaService.getInstance()
+ .updateSemanticModel(
+ original.nameIdentifier(),
+ ignored -> {
+ SemanticModelMetaService.getInstance()
+ .deleteSemanticModel(original.nameIdentifier());
+ return expected;
+ }));
+
+ Map<Integer, VersionState> versions =
listSemanticModelVersions(original.id());
+ assertEquals(1, versions.size());
+ assertTrue(versions.containsKey(1));
+ assertTrue(versions.get(1).deletedAt > 0L);
+ }
+
+ @TestTemplate
+ public void testOptimisticDropRejectsStaleInternalVersion() throws
IOException {
+ String semanticModelName =
GravitinoITUtils.genRandomName("optimistic_drop_model");
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ semanticModelName,
+ "orders_v1",
+ "v1 comment",
+ "v1");
+ SemanticModelMetaService.getInstance().insertSemanticModel(original,
false);
+ SemanticModelPO stalePO =
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelPOByIdentifier(original.nameIdentifier());
+ SemanticModelEntity updated =
+ semanticModelEntity(original.id(), semanticModelName, "orders_v2", "v2
comment", "v2");
+ SemanticModelMetaService.getInstance()
+ .updateSemanticModel(original.nameIdentifier(), ignored -> updated);
+
+ assertThrows(
+ OptimisticLockException.class,
+ () ->
+ SemanticModelMetaService.getInstance()
+ .deleteSemanticModelWithVersion(original.nameIdentifier(),
stalePO));
+
+ assertEquals(
+ updated,
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(updated.nameIdentifier()));
+ Map<Integer, VersionState> versions =
listSemanticModelVersions(original.id());
+ assertEquals(2, versions.size());
+ assertEquals(0L, versions.get(1).deletedAt);
+ assertEquals(0L, versions.get(2).deletedAt);
+ }
+
+ @TestTemplate
+ public void testStaleDropReportsAlreadyDeleted() throws IOException {
+ String semanticModelName =
GravitinoITUtils.genRandomName("already_deleted_model");
+ SemanticModelEntity original =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(), semanticModelName, "orders",
"original", "v1");
+ SemanticModelMetaService service = SemanticModelMetaService.getInstance();
+ service.insertSemanticModel(original, false);
+ SemanticModelPO stalePO =
service.getSemanticModelPOByIdentifier(original.nameIdentifier());
+
+ assertTrue(service.deleteSemanticModel(original.nameIdentifier()));
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
service.deleteSemanticModelWithVersion(original.nameIdentifier(), stalePO));
+ assertEquals(1, listSemanticModelVersions(original.id()).size());
+ }
+
+ @TestTemplate
+ public void testSoftDropVersionRetentionAndLegacyGarbageCollection() throws
IOException {
+ String semanticModelName =
GravitinoITUtils.genRandomName("retention_model");
+ SemanticModelEntity current =
+ semanticModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ semanticModelName,
+ "orders_v1",
+ "v1 comment",
+ "v1");
+ SemanticModelMetaService.getInstance().insertSemanticModel(current, false);
+ for (int version = 2; version <= 3; version++) {
+ SemanticModelEntity next =
+ semanticModelEntity(
+ current.id(),
+ semanticModelName,
+ "orders_v" + version,
+ "v" + version + " comment",
+ "v" + version);
+ current =
+ SemanticModelMetaService.getInstance()
+ .updateSemanticModel(current.nameIdentifier(), ignored -> next);
+ }
+
+ assertEquals(
+ 2,
+ SemanticModelMetaService.getInstance()
+ .deleteSemanticModelVersionsByRetentionCount(1L, 100));
+ Map<Integer, VersionState> retainedVersions =
listSemanticModelVersions(current.id());
+ assertTrue(retainedVersions.get(1).deletedAt > 0L);
+ assertTrue(retainedVersions.get(2).deletedAt > 0L);
+ assertEquals(0L, retainedVersions.get(3).deletedAt);
+
+ SemanticModelEntity finalCurrent = current;
+ long lastChangeId = maxEntityChangeId();
+ assertTrue(
+
SemanticModelMetaService.getInstance().deleteSemanticModel(finalCurrent.nameIdentifier()));
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ SemanticModelMetaService.getInstance()
+ .getSemanticModelByIdentifier(finalCurrent.nameIdentifier()));
+ assertTrue(listSemanticModelVersions(finalCurrent.id()).get(3).deletedAt >
0L);
+ assertEntityChange(lastChangeId, semanticModelName, OperateType.DROP);
+
+ int deleted =
+ SemanticModelMetaService.getInstance()
+
.deleteSemanticModelMetasByLegacyTimeline(Instant.now().toEpochMilli() + 1000,
100);
+ assertEquals(4, deleted);
+ assertEquals(0, listSemanticModelVersions(finalCurrent.id()).size());
+ assertEquals(0, countSemanticModelIdentities(finalCurrent.id()));
+ }
+
+ private SemanticModelEntity semanticModelEntity(
+ Long id, String name, String datasetName, String comment, String
extensionData) {
+ return SemanticModelEntity.builder()
+ .withId(id)
+ .withName(name)
+ .withNamespace(namespace)
+ .withComment(comment)
+ .withDefinition(
+ SemanticModelDefinition.builder()
+ .withAIContext(AIContext.of("AI context for " + datasetName))
+ .withDatasets(
+ new Dataset[] {
+ Dataset.builder()
+ .withName(datasetName)
+ .withSource(NameIdentifier.of("sales", "mart",
datasetName))
+ .build()
+ })
+ .withCustomExtensions(
+ new CustomExtension[] {
+ CustomExtension.builder()
+ .withVendorName("test")
+ .withData(extensionData)
+ .build()
+ })
+ .build())
+ .withProperties(ImmutableMap.of("dataset", datasetName))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ }
+
+ private SchemaPO getObservedSchemaPO() {
+ return SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class, mapper ->
mapper.selectSchemaMetaById(schemaId));
+ }
+
+ private Future<Throwable> startSchemaCascadeDelete(
+ SchemaPO observedSchemaPO,
+ CountDownLatch schemaDeleteLocked,
+ CountDownLatch allowDeleteCommit,
+ ExecutorService executor) {
+ return executor.submit(
+ () -> {
+ try {
+ // Reproduce the schema-row lock and Semantic Model cleanup
portion of
+ // SchemaMetaService.deleteSchema(..., true) in one controllable
transaction. The
+ // public cascade has no test hook between those operations.
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ int deleted =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+ mapper.softDeleteSchemaMetaBySchemaIdAndVersion(
+ observedSchemaPO.getSchemaId(),
+ observedSchemaPO.getCurrentVersion()));
+ assertEquals(1, deleted);
+ schemaDeleteLocked.countDown();
+ try {
+ assertTrue(allowDeleteCommit.await(30, TimeUnit.SECONDS));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelMetaMapper.class,
+ mapper ->
+ mapper.softDeleteSemanticModelMetasBySchemaIds(
+ List.of(observedSchemaPO.getSchemaId()))),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SemanticModelVersionInfoMapper.class,
+ mapper ->
+ mapper.softDeleteSemanticModelVersionsBySchemaIds(
+ List.of(observedSchemaPO.getSchemaId()))));
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ }
+
+ private Map<Integer, VersionState> listSemanticModelVersions(Long
semanticModelId) {
+ Map<Integer, VersionState> versions = new HashMap<>();
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet resultSet =
+ statement.executeQuery(
+ String.format(
+ "SELECT version, semantic_model_name,
semantic_model_comment,"
+ + " semantic_model_definition, deleted_at"
+ + " FROM semantic_model_version_info WHERE
semantic_model_id = %d",
+ semanticModelId))) {
+ while (resultSet.next()) {
+ versions.put(
+ resultSet.getInt("version"),
+ new VersionState(
+ resultSet.getString("semantic_model_name"),
+ resultSet.getString("semantic_model_comment"),
+ resultSet.getString("semantic_model_definition"),
+ resultSet.getLong("deleted_at")));
+ }
+ return versions;
+ } catch (SQLException e) {
+ throw new RuntimeException("SQL execution failed", e);
+ }
+ }
+
+ private long[] semanticModelIdentityVersions(Long semanticModelId) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet resultSet =
+ statement.executeQuery(
+ String.format(
+ "SELECT current_version, last_version FROM
semantic_model_meta"
+ + " WHERE semantic_model_id = %d",
+ semanticModelId))) {
+ assertTrue(resultSet.next());
+ return new long[] {resultSet.getLong("current_version"),
resultSet.getLong("last_version")};
+ } catch (SQLException e) {
+ throw new RuntimeException("SQL execution failed", e);
+ }
+ }
+
+ private int countSemanticModelIdentities(Long semanticModelId) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet resultSet =
+ statement.executeQuery(
+ String.format(
+ "SELECT count(*) FROM semantic_model_meta WHERE
semantic_model_id = %d",
+ semanticModelId))) {
+ assertTrue(resultSet.next());
+ return resultSet.getInt(1);
+ } catch (SQLException e) {
+ throw new RuntimeException("SQL execution failed", e);
+ }
+ }
+
+ private long maxEntityChangeId() {
+ return SessionUtils.doWithCommitAndFetchResult(
+ EntityChangeLogMapper.class, EntityChangeLogMapper::selectMaxChangeId);
+ }
+
+ private void assertEntityChange(
+ long lastConsumedId, String semanticModelName, OperateType operateType) {
+ List<EntityChangeRecord> changes =
+ SessionUtils.doWithCommitAndFetchResult(
+ EntityChangeLogMapper.class, mapper ->
mapper.selectEntityChanges(lastConsumedId, 100));
+ String fullName =
+ EntityChangeLogNameIdentifierCodec.encode(
+ NameIdentifierUtil.ofSemanticModel(
+ metalakeName, catalogName, schemaName, semanticModelName));
+ assertTrue(
+ changes.stream()
+ .anyMatch(
+ record ->
+ record.getMetalakeName().equals(metalakeName)
+ &&
record.getEntityType().equals(Entity.EntityType.SEMANTIC_MODEL.name())
+ && record.getFullName().equals(fullName)
+ && record.getOperateType() == operateType));
+ }
+
+ private static class VersionState {
+ private final String name;
+ private final String comment;
+ private final String definition;
+ private final long deletedAt;
+
+ private VersionState(String name, String comment, String definition, long
deletedAt) {
+ this.name = name;
+ this.comment = comment;
+ this.definition = definition;
+ this.deletedAt = deletedAt;
+ }
+ }
+}