This is an automated email from the ASF dual-hosted git repository.
jerryshao 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 415906bab2 [#12653] improvement(core): add OCC for view writes (#12824)
415906bab2 is described below
commit 415906bab2006d49d669047ba38daa4b423b8017
Author: Qi Yu <[email protected]>
AuthorDate: Tue Sep 8 16:36:41 2026 +0800
[#12653] improvement(core): add OCC for view writes (#12824)
### What changes were proposed in this pull request?
- Add version-CAS updates and deletes for managed views.
- Execute root metadata CAS before writing view versions and dependent
rows.
- Classify stale-version conflicts and identity changes through
stable-ID locking reads.
- Use a dedicated root-only locking query for overwrite while keeping
normal reads on the strict current-version INNER JOIN.
- Preserve stored view IDs and monotonic versions for natural-key
overwrite on H2, MySQL, and PostgreSQL.
- Fence parent schemas during create and the correct target ancestry
during cross-catalog moves.
- Add real two-transaction coverage for parent deletion,
overwrite/rename, and target-schema deletion races.
### Why are the changes needed?
Concurrent writers could create competing view-version rows, stale
deletes could remove relationships belonging to newer views, and
overwrite or move operations could violate identity, version, or parent
invariants.
Fix: #12653
### Does this PR introduce _any_ user-facing change?
No API or configuration changes. Concurrent view writes now report
optimistic-lock or not-found errors consistently.
### How was this patch tested?
- `./gradlew :core:spotlessApply :core:compileTestJava :core:javadoc`
- Ran `TestViewMetaService`, `TestViewMetaBaseSQLProvider`, and
`TestViewMetaPostgreSQLProvider`.
- Service tests ran against H2, MySQL, and PostgreSQL.
---
.../storage/relational/mapper/TagMetaMapper.java | 10 +
.../storage/relational/mapper/ViewMetaMapper.java | 46 +-
.../mapper/ViewMetaSQLProviderFactory.java | 32 +-
.../relational/mapper/ViewVersionInfoMapper.java | 11 +-
.../mapper/ViewVersionInfoSQLProviderFactory.java | 10 +-
.../provider/base/ViewMetaBaseSQLProvider.java | 79 +-
.../base/ViewVersionInfoBaseSQLProvider.java | 51 +-
.../postgresql/ViewMetaPostgreSQLProvider.java | 36 +-
.../ViewVersionInfoPostgreSQLProvider.java | 51 +-
.../relational/service/CatalogMetaService.java | 8 +-
.../relational/service/MetalakeMetaService.java | 7 +-
.../relational/service/OccWriteSupport.java | 29 +
.../relational/service/SchemaMetaService.java | 7 +-
.../storage/relational/service/TagMetaService.java | 52 +-
.../relational/service/ViewMetaService.java | 300 ++++++--
.../relational/service/ViewPOStorageOps.java | 17 +-
.../provider/base/TestViewMetaBaseSQLProvider.java | 76 ++
.../postgresql/TestViewMetaPostgreSQLProvider.java | 41 ++
.../relational/service/TestOccWriteSupport.java | 45 ++
.../relational/service/TestTagMetaService.java | 114 +++
.../relational/service/TestViewMetaService.java | 805 ++++++++++++++++++++-
21 files changed, 1570 insertions(+), 257 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetaMapper.java
index 721499826f..246e6b03bc 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetaMapper.java
@@ -23,6 +23,7 @@ import org.apache.gravitino.storage.relational.po.TagPO;
import org.apache.ibatis.annotations.DeleteProvider;
import org.apache.ibatis.annotations.InsertProvider;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.annotations.Select;
import org.apache.ibatis.annotations.SelectProvider;
import org.apache.ibatis.annotations.UpdateProvider;
@@ -30,6 +31,15 @@ public interface TagMetaMapper {
String TAG_TABLE_NAME = "tag_meta";
+ /**
+ * Counts deleted rows that still reserve the requested tag ID.
+ *
+ * @param tagId The tag ID.
+ * @return The number of deleted rows with this ID.
+ */
+ @Select("SELECT COUNT(*) FROM " + TAG_TABLE_NAME + " WHERE tag_id = #{tagId}
AND deleted_at > 0")
+ int countDeletedTagMetasById(@Param("tagId") Long tagId);
+
@SelectProvider(type = TagMetaSQLProviderFactory.class, method =
"listTagPOsByMetalake")
List<TagPO> listTagPOsByMetalake(@Param("metalakeName") String metalakeName);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java
index 813d646c15..a2d3bcf125 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java
@@ -107,6 +107,37 @@ public interface ViewMetaMapper {
ViewPO selectViewMetaBySchemaIdAndName(
@Param("schemaId") Long schemaId, @Param("viewName") String name);
+ /**
+ * Selects and exclusively locks an active view by its natural key without
joining a version row.
+ *
+ * @param schemaId the schema ID
+ * @param name the view name
+ * @return the locked view root, or {@code null} when it does not exist
+ */
+ @SelectProvider(
+ type = ViewMetaSQLProviderFactory.class,
+ method = "selectViewMetaBySchemaIdAndNameForUpdate")
+ ViewPO selectViewMetaBySchemaIdAndNameForUpdate(
+ @Param("schemaId") Long schemaId, @Param("viewName") String name);
+
+ /**
+ * Selects and exclusively locks an active view metadata row.
+ *
+ * @param viewId the view ID
+ * @return the active view metadata, or {@code null} when it no longer exists
+ */
+ @SelectProvider(type = ViewMetaSQLProviderFactory.class, method =
"selectViewMetaByIdForUpdate")
+ ViewPO selectViewMetaByIdForUpdate(@Param("viewId") Long viewId);
+
+ /**
+ * Checks whether a soft-deleted view still owns the requested primary key.
+ *
+ * @param viewId the view ID
+ * @return one if a deleted row reserves the ID, otherwise zero
+ */
+ @Select("SELECT COUNT(*) FROM " + TABLE_NAME + " WHERE view_id = #{viewId}
AND deleted_at > 0")
+ int countDeletedViewMetasById(@Param("viewId") Long viewId);
+
@ResultMap("viewPOResultMap")
@SelectProvider(type = ViewMetaSQLProviderFactory.class, method =
"selectViewByFullQualifiedName")
ViewPO selectViewByFullQualifiedName(
@@ -118,17 +149,20 @@ public interface ViewMetaMapper {
@InsertProvider(type = ViewMetaSQLProviderFactory.class, method =
"insertViewMeta")
void insertViewMeta(@Param("viewMeta") ViewPO viewPO);
- @InsertProvider(
- type = ViewMetaSQLProviderFactory.class,
- method = "insertViewMetaOnDuplicateKeyUpdate")
- void insertViewMetaOnDuplicateKeyUpdate(@Param("viewMeta") ViewPO viewPO);
-
@UpdateProvider(type = ViewMetaSQLProviderFactory.class, method =
"updateViewMeta")
Integer updateViewMeta(
@Param("newViewMeta") ViewPO newViewPO, @Param("oldViewMeta") ViewPO
oldViewPO);
+ /**
+ * Soft-deletes a view only if its version has not changed since the caller
read it.
+ *
+ * @param viewId the view ID
+ * @param currentVersion the version observed by the caller
+ * @return the number of deleted rows; zero means the view changed or
disappeared
+ */
@UpdateProvider(type = ViewMetaSQLProviderFactory.class, method =
"softDeleteViewMetasByViewId")
- Integer softDeleteViewMetasByViewId(@Param("viewId") Long viewId);
+ Integer softDeleteViewMetasByViewId(
+ @Param("viewId") Long viewId, @Param("currentVersion") Long
currentVersion);
@UpdateProvider(
type = ViewMetaSQLProviderFactory.class,
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java
index cb25d743d3..8531c2f748 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java
@@ -76,6 +76,22 @@ public class ViewMetaSQLProviderFactory {
return getProvider().selectViewMetaBySchemaIdAndName(schemaId, name);
}
+ /** Delegates a root-only locking lookup by view natural key. */
+ public static String selectViewMetaBySchemaIdAndNameForUpdate(
+ @Param("schemaId") Long schemaId, @Param("viewName") String name) {
+ return getProvider().selectViewMetaBySchemaIdAndNameForUpdate(schemaId,
name);
+ }
+
+ /**
+ * Returns SQL that selects and exclusively locks an active view metadata
row.
+ *
+ * @param viewId the view ID
+ * @return the locking select SQL
+ */
+ public static String selectViewMetaByIdForUpdate(@Param("viewId") Long
viewId) {
+ return getProvider().selectViewMetaByIdForUpdate(viewId);
+ }
+
public static String selectViewByFullQualifiedName(
@Param("metalakeName") String metalakeName,
@Param("catalogName") String catalogName,
@@ -89,17 +105,21 @@ public class ViewMetaSQLProviderFactory {
return getProvider().insertViewMeta(viewPO);
}
- public static String insertViewMetaOnDuplicateKeyUpdate(@Param("viewMeta")
ViewPO viewPO) {
- return getProvider().insertViewMetaOnDuplicateKeyUpdate(viewPO);
- }
-
public static String updateViewMeta(
@Param("newViewMeta") ViewPO newViewPO, @Param("oldViewMeta") ViewPO
oldViewPO) {
return getProvider().updateViewMeta(newViewPO, oldViewPO);
}
- public static String softDeleteViewMetasByViewId(@Param("viewId") Long
viewId) {
- return getProvider().softDeleteViewMetasByViewId(viewId);
+ /**
+ * Returns SQL that soft-deletes a view with a version check.
+ *
+ * @param viewId the view ID
+ * @param currentVersion the version observed by the caller
+ * @return the version-checked delete SQL
+ */
+ public static String softDeleteViewMetasByViewId(
+ @Param("viewId") Long viewId, @Param("currentVersion") Long
currentVersion) {
+ return getProvider().softDeleteViewMetasByViewId(viewId, currentVersion);
}
public static String softDeleteViewMetasByMetalakeId(@Param("metalakeId")
Long metalakeId) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoMapper.java
index 65cf74e83a..c64622d551 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoMapper.java
@@ -18,6 +18,7 @@
*/
package org.apache.gravitino.storage.relational.mapper;
+import java.util.List;
import org.apache.gravitino.storage.relational.po.ViewVersionInfoPO;
import org.apache.ibatis.annotations.DeleteProvider;
import org.apache.ibatis.annotations.InsertProvider;
@@ -33,12 +34,6 @@ public interface ViewVersionInfoMapper {
@InsertProvider(type = ViewVersionInfoSQLProviderFactory.class, method =
"insertViewVersionInfo")
void insertViewVersionInfo(@Param("viewVersionInfo") ViewVersionInfoPO
viewVersionInfoPO);
- @InsertProvider(
- type = ViewVersionInfoSQLProviderFactory.class,
- method = "insertViewVersionInfoOnDuplicateKeyUpdate")
- void insertViewVersionInfoOnDuplicateKeyUpdate(
- @Param("viewVersionInfo") ViewVersionInfoPO viewVersionInfoPO);
-
@SelectProvider(
type = ViewVersionInfoSQLProviderFactory.class,
method = "selectViewVersionInfoByViewIdAndVersion")
@@ -52,8 +47,8 @@ public interface ViewVersionInfoMapper {
@UpdateProvider(
type = ViewVersionInfoSQLProviderFactory.class,
- method = "softDeleteViewVersionsBySchemaId")
- Integer softDeleteViewVersionsBySchemaId(@Param("schemaId") Long schemaId);
+ method = "softDeleteViewVersionsBySchemaIds")
+ Integer softDeleteViewVersionsBySchemaIds(@Param("schemaIds") List<Long>
schemaIds);
@UpdateProvider(
type = ViewVersionInfoSQLProviderFactory.class,
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoSQLProviderFactory.java
index e04d0237d3..5df1190a3c 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewVersionInfoSQLProviderFactory.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.ViewVersionInfoBaseSQLProvider;
@@ -56,11 +57,6 @@ public class ViewVersionInfoSQLProviderFactory {
return getProvider().insertViewVersionInfo(viewVersionInfoPO);
}
- public static String insertViewVersionInfoOnDuplicateKeyUpdate(
- @Param("viewVersionInfo") ViewVersionInfoPO viewVersionInfoPO) {
- return
getProvider().insertViewVersionInfoOnDuplicateKeyUpdate(viewVersionInfoPO);
- }
-
public static String selectViewVersionInfoByViewIdAndVersion(
@Param("viewId") Long viewId, @Param("version") Integer version) {
return getProvider().selectViewVersionInfoByViewIdAndVersion(viewId,
version);
@@ -70,8 +66,8 @@ public class ViewVersionInfoSQLProviderFactory {
return getProvider().softDeleteViewVersionsByViewId(viewId);
}
- public static String softDeleteViewVersionsBySchemaId(@Param("schemaId")
Long schemaId) {
- return getProvider().softDeleteViewVersionsBySchemaId(schemaId);
+ public static String softDeleteViewVersionsBySchemaIds(@Param("schemaIds")
List<Long> schemaIds) {
+ return getProvider().softDeleteViewVersionsBySchemaIds(schemaIds);
}
public static String softDeleteViewVersionsByCatalogId(@Param("catalogId")
Long catalogId) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java
index 1eaf2eabef..4f5ecbb04a 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java
@@ -131,26 +131,45 @@ public class ViewMetaBaseSQLProvider {
+ " AND vm.deleted_at = 0 AND vi.deleted_at = 0";
}
- public String insertViewMeta(@Param("viewMeta") ViewPO viewPO) {
- return "INSERT INTO "
+ /**
+ * Returns SQL that locks an active view by natural key without joining its
version row.
+ *
+ * <p>This query is reserved for overwrite decisions. Normal reads keep
using the inner-joined
+ * query above so a broken current-version invariant is reported as missing
instead of producing a
+ * partially populated {@code ViewPO}.
+ */
+ public String selectViewMetaBySchemaIdAndNameForUpdate(
+ @Param("schemaId") Long schemaId, @Param("viewName") String name) {
+ return "SELECT view_id as viewId, view_name as viewName,"
+ + " metalake_id as metalakeId, catalog_id as catalogId, schema_id as
schemaId,"
+ + " current_version as currentVersion, last_version as lastVersion,"
+ + " deleted_at as deletedAt"
+ + " FROM "
+ TABLE_NAME
- + " (view_id, view_name, metalake_id,"
- + " catalog_id, schema_id,"
- + " current_version, last_version, audit_info, deleted_at)"
- + " VALUES ("
- + " #{viewMeta.viewId},"
- + " #{viewMeta.viewName},"
- + " #{viewMeta.metalakeId},"
- + " #{viewMeta.catalogId},"
- + " #{viewMeta.schemaId},"
- + " #{viewMeta.currentVersion},"
- + " #{viewMeta.lastVersion},"
- + " #{viewMeta.auditInfo},"
- + " #{viewMeta.deletedAt}"
- + " )";
+ + " WHERE schema_id = #{schemaId} AND view_name = #{viewName}"
+ + " AND deleted_at = 0 FOR UPDATE";
}
- public String insertViewMetaOnDuplicateKeyUpdate(@Param("viewMeta") ViewPO
viewPO) {
+ /**
+ * Returns the active view metadata row and holds it exclusively for the
transaction.
+ *
+ * <p>The version table is deliberately not joined: PostgreSQL rejects
locking the nullable side
+ * of an outer join, and conflict classification only needs the root row's
identity and version.
+ *
+ * @param viewId the view ID
+ * @return the locking select SQL
+ */
+ public String selectViewMetaByIdForUpdate(@Param("viewId") Long viewId) {
+ return "SELECT view_id as viewId, view_name as viewName,"
+ + " metalake_id as metalakeId, catalog_id as catalogId, schema_id as
schemaId,"
+ + " current_version as currentVersion, last_version as lastVersion,"
+ + " audit_info as auditInfo, deleted_at as deletedAt"
+ + " FROM "
+ + TABLE_NAME
+ + " WHERE view_id = #{viewId} AND deleted_at = 0 FOR UPDATE";
+ }
+
+ public String insertViewMeta(@Param("viewMeta") ViewPO viewPO) {
return "INSERT INTO "
+ TABLE_NAME
+ " (view_id, view_name, metalake_id,"
@@ -166,16 +185,7 @@ public class ViewMetaBaseSQLProvider {
+ " #{viewMeta.lastVersion},"
+ " #{viewMeta.auditInfo},"
+ " #{viewMeta.deletedAt}"
- + " )"
- + " ON DUPLICATE KEY UPDATE"
- + " view_name = #{viewMeta.viewName},"
- + " metalake_id = #{viewMeta.metalakeId},"
- + " catalog_id = #{viewMeta.catalogId},"
- + " schema_id = #{viewMeta.schemaId},"
- + " current_version = #{viewMeta.currentVersion},"
- + " last_version = #{viewMeta.lastVersion},"
- + " audit_info = #{viewMeta.auditInfo},"
- + " deleted_at = #{viewMeta.deletedAt}";
+ + " )";
}
public String updateViewMeta(
@@ -183,6 +193,8 @@ public class ViewMetaBaseSQLProvider {
return "UPDATE "
+ TABLE_NAME
+ " SET view_name = #{newViewMeta.viewName}, "
+ + " metalake_id = #{newViewMeta.metalakeId}, "
+ + " catalog_id = #{newViewMeta.catalogId}, "
+ " schema_id = #{newViewMeta.schemaId}, "
+ " current_version = #{newViewMeta.currentVersion}, "
+ " last_version = #{newViewMeta.lastVersion}, "
@@ -215,12 +227,21 @@ public class ViewMetaBaseSQLProvider {
+ "</script>";
}
- public String softDeleteViewMetasByViewId(@Param("viewId") Long viewId) {
+ /**
+ * Returns SQL that deletes only the view version observed by the caller.
+ *
+ * @param viewId the view ID
+ * @param currentVersion the version observed by the caller
+ * @return the version-checked delete SQL
+ */
+ public String softDeleteViewMetasByViewId(
+ @Param("viewId") Long viewId, @Param("currentVersion") Long
currentVersion) {
return "UPDATE "
+ TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.MYSQL
- + " WHERE view_id = #{viewId} AND deleted_at = 0";
+ + " WHERE view_id = #{viewId}"
+ + " AND current_version = #{currentVersion} AND deleted_at = 0";
}
public String softDeleteViewMetasByMetalakeId(@Param("metalakeId") Long
metalakeId) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewVersionInfoBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewVersionInfoBaseSQLProvider.java
index d9d1010c22..d3521749ad 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewVersionInfoBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewVersionInfoBaseSQLProvider.java
@@ -18,6 +18,8 @@
*/
package org.apache.gravitino.storage.relational.mapper.provider.base;
+import java.util.List;
+import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewVersionInfoMapper;
import org.apache.gravitino.storage.relational.mapper.provider.DatabaseTimeSQL;
import org.apache.gravitino.storage.relational.po.ViewVersionInfoPO;
@@ -41,31 +43,6 @@ public class ViewVersionInfoBaseSQLProvider {
+ " #{viewVersionInfo.auditInfo}, #{viewVersionInfo.deletedAt})";
}
- public String insertViewVersionInfoOnDuplicateKeyUpdate(
- @Param("viewVersionInfo") ViewVersionInfoPO viewVersionInfoPO) {
- return "INSERT INTO "
- + ViewVersionInfoMapper.TABLE_NAME
- + " (metalake_id, catalog_id, schema_id, view_id, version,"
- + " view_comment, columns, properties, default_catalog,
default_schema, representations,"
- + " audit_info, deleted_at)"
- + " VALUES (#{viewVersionInfo.metalakeId},
#{viewVersionInfo.catalogId},"
- + " #{viewVersionInfo.schemaId}, #{viewVersionInfo.viewId},"
- + " #{viewVersionInfo.version}, #{viewVersionInfo.viewComment},"
- + " #{viewVersionInfo.columns}, #{viewVersionInfo.properties},"
- + " #{viewVersionInfo.defaultCatalog},
#{viewVersionInfo.defaultSchema},"
- + " #{viewVersionInfo.representations},"
- + " #{viewVersionInfo.auditInfo}, #{viewVersionInfo.deletedAt})"
- + " ON DUPLICATE KEY UPDATE"
- + " view_comment = #{viewVersionInfo.viewComment},"
- + " columns = #{viewVersionInfo.columns},"
- + " properties = #{viewVersionInfo.properties},"
- + " default_catalog = #{viewVersionInfo.defaultCatalog},"
- + " default_schema = #{viewVersionInfo.defaultSchema},"
- + " representations = #{viewVersionInfo.representations},"
- + " audit_info = #{viewVersionInfo.auditInfo},"
- + " deleted_at = #{viewVersionInfo.deletedAt}";
- }
-
public String selectViewVersionInfoByViewIdAndVersion(
@Param("viewId") Long viewId, @Param("version") Integer version) {
return "SELECT id as id, metalake_id as metalakeId, catalog_id as
catalogId,"
@@ -86,12 +63,22 @@ public class ViewVersionInfoBaseSQLProvider {
+ " WHERE view_id = #{viewId} AND deleted_at = 0";
}
- public String softDeleteViewVersionsBySchemaId(@Param("schemaId") Long
schemaId) {
- return "UPDATE "
+ public String softDeleteViewVersionsBySchemaIds(@Param("schemaIds")
List<Long> schemaIds) {
+ return "<script>"
+ + "UPDATE "
+ ViewVersionInfoMapper.TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.MYSQL
- + " WHERE schema_id = #{schemaId} AND deleted_at = 0";
+ // History follows the stable entity ID, not the parent recorded in
each snapshot.
+ // Include deleted roots: cascade cleanup soft-deletes roots before
their versions.
+ + " WHERE view_id IN (SELECT view_id FROM "
+ + ViewMetaMapper.TABLE_NAME
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}"
+ + "</foreach>"
+ + ")) AND deleted_at = 0"
+ + "</script>";
}
public String softDeleteViewVersionsByCatalogId(@Param("catalogId") Long
catalogId) {
@@ -99,7 +86,9 @@ public class ViewVersionInfoBaseSQLProvider {
+ ViewVersionInfoMapper.TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.MYSQL
- + " WHERE catalog_id = #{catalogId} AND deleted_at = 0";
+ + " WHERE view_id IN (SELECT view_id FROM "
+ + ViewMetaMapper.TABLE_NAME
+ + " WHERE catalog_id = #{catalogId}) AND deleted_at = 0";
}
public String softDeleteViewVersionsByMetalakeId(@Param("metalakeId") Long
metalakeId) {
@@ -107,7 +96,9 @@ public class ViewVersionInfoBaseSQLProvider {
+ ViewVersionInfoMapper.TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.MYSQL
- + " WHERE metalake_id = #{metalakeId} AND deleted_at = 0";
+ + " WHERE view_id IN (SELECT view_id FROM "
+ + ViewMetaMapper.TABLE_NAME
+ + " WHERE metalake_id = #{metalakeId}) AND deleted_at = 0";
}
public String deleteViewVersionsByLegacyTimeline(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java
index 343375c696..058e431736 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java
@@ -24,40 +24,10 @@ import static
org.apache.gravitino.storage.relational.mapper.ViewMetaMapper.VERS
import java.util.List;
import org.apache.gravitino.storage.relational.mapper.provider.DatabaseTimeSQL;
import
org.apache.gravitino.storage.relational.mapper.provider.base.ViewMetaBaseSQLProvider;
-import org.apache.gravitino.storage.relational.po.ViewPO;
import org.apache.ibatis.annotations.Param;
public class ViewMetaPostgreSQLProvider extends ViewMetaBaseSQLProvider {
- @Override
- public String insertViewMetaOnDuplicateKeyUpdate(@Param("viewMeta") ViewPO
viewPO) {
- return "INSERT INTO "
- + TABLE_NAME
- + " (view_id, view_name, metalake_id,"
- + " catalog_id, schema_id,"
- + " current_version, last_version, audit_info, deleted_at)"
- + " VALUES ("
- + " #{viewMeta.viewId},"
- + " #{viewMeta.viewName},"
- + " #{viewMeta.metalakeId},"
- + " #{viewMeta.catalogId},"
- + " #{viewMeta.schemaId},"
- + " #{viewMeta.currentVersion},"
- + " #{viewMeta.lastVersion},"
- + " #{viewMeta.auditInfo},"
- + " #{viewMeta.deletedAt}"
- + " )"
- + " ON CONFLICT (view_id) DO UPDATE SET"
- + " view_name = #{viewMeta.viewName},"
- + " metalake_id = #{viewMeta.metalakeId},"
- + " catalog_id = #{viewMeta.catalogId},"
- + " schema_id = #{viewMeta.schemaId},"
- + " current_version = #{viewMeta.currentVersion},"
- + " last_version = #{viewMeta.lastVersion},"
- + " audit_info = #{viewMeta.auditInfo},"
- + " deleted_at = #{viewMeta.deletedAt}";
- }
-
@Override
public String listViewPOsBySchemaId(@Param("schemaId") Long schemaId) {
return "SELECT vm.view_id, vm.view_name, vm.metalake_id, vm.catalog_id,
vm.schema_id,"
@@ -95,12 +65,14 @@ public class ViewMetaPostgreSQLProvider extends
ViewMetaBaseSQLProvider {
}
@Override
- public String softDeleteViewMetasByViewId(@Param("viewId") Long viewId) {
+ public String softDeleteViewMetasByViewId(
+ @Param("viewId") Long viewId, @Param("currentVersion") Long
currentVersion) {
return "UPDATE "
+ TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.POSTGRESQL
- + " WHERE view_id = #{viewId} AND deleted_at = 0";
+ + " WHERE view_id = #{viewId}"
+ + " AND current_version = #{currentVersion} AND deleted_at = 0";
}
@Override
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewVersionInfoPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewVersionInfoPostgreSQLProvider.java
index 0fd96bc92b..3b0a8307ee 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewVersionInfoPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewVersionInfoPostgreSQLProvider.java
@@ -18,40 +18,15 @@
*/
package org.apache.gravitino.storage.relational.mapper.provider.postgresql;
+import java.util.List;
+import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewVersionInfoMapper;
import org.apache.gravitino.storage.relational.mapper.provider.DatabaseTimeSQL;
import
org.apache.gravitino.storage.relational.mapper.provider.base.ViewVersionInfoBaseSQLProvider;
-import org.apache.gravitino.storage.relational.po.ViewVersionInfoPO;
import org.apache.ibatis.annotations.Param;
public class ViewVersionInfoPostgreSQLProvider extends
ViewVersionInfoBaseSQLProvider {
- @Override
- public String insertViewVersionInfoOnDuplicateKeyUpdate(
- @Param("viewVersionInfo") ViewVersionInfoPO viewVersionInfoPO) {
- return "INSERT INTO "
- + ViewVersionInfoMapper.TABLE_NAME
- + " (metalake_id, catalog_id, schema_id, view_id, version,"
- + " view_comment, columns, properties, default_catalog,
default_schema, representations,"
- + " audit_info, deleted_at)"
- + " VALUES (#{viewVersionInfo.metalakeId},
#{viewVersionInfo.catalogId},"
- + " #{viewVersionInfo.schemaId}, #{viewVersionInfo.viewId},"
- + " #{viewVersionInfo.version}, #{viewVersionInfo.viewComment},"
- + " #{viewVersionInfo.columns}, #{viewVersionInfo.properties},"
- + " #{viewVersionInfo.defaultCatalog},
#{viewVersionInfo.defaultSchema},"
- + " #{viewVersionInfo.representations},"
- + " #{viewVersionInfo.auditInfo}, #{viewVersionInfo.deletedAt})"
- + " ON CONFLICT (view_id, version, deleted_at) DO UPDATE SET"
- + " view_comment = #{viewVersionInfo.viewComment},"
- + " columns = #{viewVersionInfo.columns},"
- + " properties = #{viewVersionInfo.properties},"
- + " default_catalog = #{viewVersionInfo.defaultCatalog},"
- + " default_schema = #{viewVersionInfo.defaultSchema},"
- + " representations = #{viewVersionInfo.representations},"
- + " audit_info = #{viewVersionInfo.auditInfo},"
- + " deleted_at = #{viewVersionInfo.deletedAt}";
- }
-
@Override
public String softDeleteViewVersionsByViewId(@Param("viewId") Long viewId) {
return "UPDATE "
@@ -62,12 +37,20 @@ public class ViewVersionInfoPostgreSQLProvider extends
ViewVersionInfoBaseSQLPro
}
@Override
- public String softDeleteViewVersionsBySchemaId(@Param("schemaId") Long
schemaId) {
- return "UPDATE "
+ public String softDeleteViewVersionsBySchemaIds(@Param("schemaIds")
List<Long> schemaIds) {
+ return "<script>"
+ + "UPDATE "
+ ViewVersionInfoMapper.TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.POSTGRESQL
- + " WHERE schema_id = #{schemaId} AND deleted_at = 0";
+ + " WHERE view_id IN (SELECT view_id FROM "
+ + ViewMetaMapper.TABLE_NAME
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}"
+ + "</foreach>"
+ + ")) AND deleted_at = 0"
+ + "</script>";
}
@Override
@@ -76,7 +59,9 @@ public class ViewVersionInfoPostgreSQLProvider extends
ViewVersionInfoBaseSQLPro
+ ViewVersionInfoMapper.TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.POSTGRESQL
- + " WHERE catalog_id = #{catalogId} AND deleted_at = 0";
+ + " WHERE view_id IN (SELECT view_id FROM "
+ + ViewMetaMapper.TABLE_NAME
+ + " WHERE catalog_id = #{catalogId}) AND deleted_at = 0";
}
@Override
@@ -85,7 +70,9 @@ public class ViewVersionInfoPostgreSQLProvider extends
ViewVersionInfoBaseSQLPro
+ ViewVersionInfoMapper.TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.POSTGRESQL
- + " WHERE metalake_id = #{metalakeId} AND deleted_at = 0";
+ + " WHERE view_id IN (SELECT view_id FROM "
+ + ViewMetaMapper.TABLE_NAME
+ + " WHERE metalake_id = #{metalakeId}) AND deleted_at = 0";
}
@Override
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 66970fd4cc..d7895ec7c8 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
@@ -55,6 +55,7 @@ import
org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ViewVersionInfoMapper;
import org.apache.gravitino.storage.relational.po.CatalogPO;
import org.apache.gravitino.storage.relational.po.MetalakePO;
import org.apache.gravitino.storage.relational.po.SchemaPO;
@@ -347,8 +348,11 @@ public class CatalogMetaService {
mapper -> mapper.softDeleteStatisticsByCatalogId(catalogId)),
() ->
SessionUtils.doWithoutCommit(
- ViewMetaMapper.class,
- mapper -> mapper.softDeleteViewMetasByCatalogId(catalogId)));
+ ViewMetaMapper.class, mapper ->
mapper.softDeleteViewMetasByCatalogId(catalogId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.softDeleteViewVersionsByCatalogId(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 e7cf7a9cf5..174cd97df1 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
@@ -63,6 +63,7 @@ import
org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.UserMetaMapper;
import org.apache.gravitino.storage.relational.mapper.UserRoleRelMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ViewVersionInfoMapper;
import org.apache.gravitino.storage.relational.po.CatalogPO;
import org.apache.gravitino.storage.relational.po.MetalakePO;
import org.apache.gravitino.storage.relational.po.SchemaPO;
@@ -326,7 +327,11 @@ public class MetalakeMetaService {
() ->
SessionUtils.doWithoutCommit(
ViewMetaMapper.class,
- mapper ->
mapper.softDeleteViewMetasByMetalakeId(metalakeId)));
+ mapper ->
mapper.softDeleteViewMetasByMetalakeId(metalakeId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.softDeleteViewVersionsByMetalakeId(metalakeId)));
} else {
SessionUtils.doMultipleWithCommit(
() -> {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/OccWriteSupport.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/OccWriteSupport.java
index 01842bfe19..c1cc7af577 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/OccWriteSupport.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/OccWriteSupport.java
@@ -27,6 +27,7 @@ import java.util.function.Supplier;
import java.util.function.ToIntFunction;
import javax.annotation.Nullable;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.storage.relational.utils.ExceptionUtils;
@@ -39,6 +40,34 @@ public class OccWriteSupport {
private OccWriteSupport() {}
+ /**
+ * Resolves an overwrite target by locking its natural key, then falling
back to its stable ID.
+ *
+ * <p>A missing row is not locked. Callers must handle concurrent insert
conflicts and roll back
+ * the whole write before retrying. A stable ID under another parent must
never be adopted without
+ * fencing that parent.
+ *
+ * @param <T> the persistent object type
+ * @param byNameLookup the locking natural-key lookup
+ * @param byIdLookup the locking stable-ID lookup
+ * @param sameParent checks whether the stable-ID match belongs to the
target parent
+ * @return the locked target, or null if both lookups miss
+ * @throws EntityAlreadyExistsException if the stable ID belongs to another
parent
+ */
+ @Nullable
+ public static <T> T findAndLockForOverwrite(
+ Supplier<T> byNameLookup, Supplier<T> byIdLookup, Predicate<T>
sameParent) {
+ T current = byNameLookup.get();
+ if (current != null) {
+ return current;
+ }
+ current = byIdLookup.get();
+ if (current != null && !sameParent.test(current)) {
+ throw new EntityAlreadyExistsException("The entity ID already belongs to
a different parent");
+ }
+ return current;
+ }
+
/**
* Classifies a write-failure for an entity during an optimistic concurrency
control operation.
*
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 037befda7e..796c0b4368 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
@@ -62,6 +62,7 @@ import
org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ViewVersionInfoMapper;
import org.apache.gravitino.storage.relational.po.CatalogPO;
import org.apache.gravitino.storage.relational.po.SchemaPO;
import org.apache.gravitino.storage.relational.utils.ExceptionUtils;
@@ -354,7 +355,11 @@ public class SchemaMetaService {
() ->
SessionUtils.doWithoutCommit(
ViewMetaMapper.class,
- mapper ->
mapper.softDeleteViewMetasBySchemaIds(schemaIds.get())));
+ mapper ->
mapper.softDeleteViewMetasBySchemaIds(schemaIds.get())),
+ () ->
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.softDeleteViewVersionsBySchemaIds(schemaIds.get())));
} else {
SessionUtils.doMultipleWithCommit(
() -> {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
index a704a74f11..01f214234b 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
@@ -124,7 +124,18 @@ public class TagMetaService {
() -> lockMetalakeForTagCreate(metalakePO),
() -> insertTagWithoutCommit(tagEntity, tagPO, overwritten));
} catch (RuntimeException e) {
- ExceptionUtils.checkSQLException(e, Entity.EntityType.TAG,
tagEntity.toString());
+ try {
+ ExceptionUtils.checkSQLException(
+ e, Entity.EntityType.TAG, tagEntity.nameIdentifier().toString());
+ } catch (EntityAlreadyExistsException duplicate) {
+ if (overwritten) {
+ // A missing-row locking read does not fence a concurrent insert at
READ_COMMITTED.
+ // Propagate the conflict so the whole transaction is rolled back
before retrying.
+ throw ExceptionUtils.concurrentModification(
+ Entity.EntityType.TAG, tagEntity.nameIdentifier());
+ }
+ throw duplicate;
+ }
throw e;
}
}
@@ -744,6 +755,14 @@ public class TagMetaService {
TagPO existingTagPO = findAndLockTagForOverwrite(initializedTagPO);
if (existingTagPO == null) {
+ if (SessionUtils.getWithoutCommit(
+ TagMetaMapper.class,
+ mapper ->
mapper.countDeletedTagMetasById(initializedTagPO.getTagId()))
+ > 0) {
+ throw new EntityAlreadyExistsException(
+ "The tag ID %s is reserved by a deleted tag; use a new ID",
+ initializedTagPO.getTagId());
+ }
insertNewTagWithoutCommit(initializedTagPO);
return;
}
@@ -759,25 +778,18 @@ public class TagMetaService {
}
private TagPO findAndLockTagForOverwrite(TagPO initializedTagPO) {
- TagPO sameNameTagPO =
- SessionUtils.getWithoutCommit(
- TagMetaMapper.class,
- mapper ->
- mapper.selectTagMetaByMetalakeIdAndNameForUpdate(
- initializedTagPO.getMetalakeId(),
initializedTagPO.getTagName()));
- if (sameNameTagPO != null) {
- return sameNameTagPO;
- }
-
- TagPO sameIdTagPO =
- SessionUtils.getWithoutCommit(
- TagMetaMapper.class,
- mapper ->
mapper.selectTagByTagIdForUpdate(initializedTagPO.getTagId()));
- if (sameIdTagPO == null
- || !Objects.equals(sameIdTagPO.getMetalakeId(),
initializedTagPO.getMetalakeId())) {
- return null;
- }
- return sameIdTagPO;
+ return OccWriteSupport.findAndLockForOverwrite(
+ () ->
+ SessionUtils.getWithoutCommit(
+ TagMetaMapper.class,
+ mapper ->
+ mapper.selectTagMetaByMetalakeIdAndNameForUpdate(
+ initializedTagPO.getMetalakeId(),
initializedTagPO.getTagName())),
+ () ->
+ SessionUtils.getWithoutCommit(
+ TagMetaMapper.class,
+ mapper ->
mapper.selectTagByTagIdForUpdate(initializedTagPO.getTagId())),
+ current -> Objects.equals(current.getMetalakeId(),
initializedTagPO.getMetalakeId()));
}
private void updateTagRootWithVersion(NameIdentifier identifier, TagPO
oldTagPO, TagPO newTagPO) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
index 2a46018446..c96f3731d1 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
@@ -27,15 +27,16 @@ import com.google.common.base.Preconditions;
import java.io.IOException;
import java.util.List;
import java.util.Objects;
-import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Function;
import java.util.stream.Collectors;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.HasIdentifier;
import org.apache.gravitino.MetadataObject;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.meta.NamespacedEntityId;
import org.apache.gravitino.meta.ViewEntity;
import org.apache.gravitino.metrics.Monitored;
import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
@@ -45,6 +46,7 @@ import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewVersionInfoMapper;
import org.apache.gravitino.storage.relational.po.ViewPO;
+import org.apache.gravitino.storage.relational.po.ViewVersionInfoPO;
import org.apache.gravitino.storage.relational.utils.ExceptionUtils;
import org.apache.gravitino.storage.relational.utils.SessionUtils;
import org.apache.gravitino.utils.NameIdentifierUtil;
@@ -116,22 +118,20 @@ public class ViewMetaService {
po.getSchemaId(),
po.getCatalogId(),
po.getMetalakeId()),
- () ->
- SessionUtils.doWithoutCommit(
- ViewMetaMapper.class, mapper -> ops.insertPO(mapper, po,
overwrite)),
- () ->
- SessionUtils.doWithoutCommit(
- ViewVersionInfoMapper.class,
- mapper -> {
- if (overwrite) {
-
mapper.insertViewVersionInfoOnDuplicateKeyUpdate(po.getViewVersionInfoPO());
- } else {
- mapper.insertViewVersionInfo(po.getViewVersionInfoPO());
- }
- }));
+ () -> insertViewWithoutCommit(viewEntity, po, overwrite));
} catch (RuntimeException re) {
- ExceptionUtils.checkSQLException(
- re, Entity.EntityType.VIEW, viewEntity.nameIdentifier().toString());
+ try {
+ ExceptionUtils.checkSQLException(
+ re, Entity.EntityType.VIEW,
viewEntity.nameIdentifier().toString());
+ } catch (EntityAlreadyExistsException duplicate) {
+ if (overwrite) {
+ // A missing-row locking read cannot fence a concurrent insert at
READ_COMMITTED.
+ // Propagate the conflict so the whole transaction is rolled back
before retrying.
+ throw ExceptionUtils.concurrentModification(
+ Entity.EntityType.VIEW, viewEntity.nameIdentifier());
+ }
+ throw duplicate;
+ }
throw re;
}
}
@@ -150,27 +150,44 @@ public class ViewMetaService {
newEntity.id(),
oldViewEntity.id());
- AtomicInteger updateResult = new AtomicInteger(0);
+ boolean isSchemaChanged =
!newEntity.namespace().equals(oldViewEntity.namespace());
+ NamespacedEntityId targetSchemaIds =
+ isSchemaChanged
+ ? EntityIdService.getEntityIds(
+ NameIdentifier.of(newEntity.namespace().levels()),
Entity.EntityType.SCHEMA)
+ : null;
+ Long newSchemaId = isSchemaChanged ? targetSchemaIds.entityId() :
oldViewPO.getSchemaId();
+ Long newCatalogId =
+ isSchemaChanged ? targetSchemaIds.namespaceIds()[1] :
oldViewPO.getCatalogId();
+ Long newMetalakeId =
+ isSchemaChanged ? targetSchemaIds.namespaceIds()[0] :
oldViewPO.getMetalakeId();
+
try {
ViewPO newViewPO = updateViewPO(oldViewPO, newEntity);
SessionUtils.doMultipleWithCommit(
- () ->
- SessionUtils.doWithoutCommit(
- ViewVersionInfoMapper.class,
- mapper ->
mapper.insertViewVersionInfo(newViewPO.getViewVersionInfoPO())),
() -> {
- updateResult.set(
+ if (isSchemaChanged) {
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ newEntity.nameIdentifier(), newSchemaId, newCatalogId,
newMetalakeId);
+ }
+ },
+ () -> {
+ // current_version is the sole OCC token. The root CAS is the
transaction's decision
+ // point and must run before the unguarded version-row insert
below.
+ int updated =
SessionUtils.getWithoutCommit(
- ViewMetaMapper.class, mapper -> ops.updatePO(mapper,
newViewPO, oldViewPO)));
- if (updateResult.get() == 0) {
- throw new RuntimeException("Failed to update the entity: " +
ident);
+ ViewMetaMapper.class, mapper -> ops.updatePO(mapper,
newViewPO, oldViewPO));
+ if (updated == 0) {
+ throw viewWriteFailure(ident, oldViewPO);
}
- });
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.insertViewVersionInfo(newViewPO.getViewVersionInfoPO())));
return newEntity;
} catch (RuntimeException re) {
- if (updateResult.get() == 0) {
- throw new IOException("Failed to update the entity: " + ident);
- }
ExceptionUtils.checkSQLException(
re, Entity.EntityType.VIEW, newEntity.nameIdentifier().toString());
throw re;
@@ -182,7 +199,30 @@ public class ViewMetaService {
baseMetricName = "deleteViewByIdentifier")
public boolean deleteView(NameIdentifier ident) {
ViewPO viewPO = getViewPOByIdentifier(ident);
- return deleteView(viewPO.getViewId());
+
+ deleteViewWithVersion(ident, viewPO);
+ return true;
+ }
+
+ /**
+ * Deletes the observed view and its dependent rows in one transaction.
+ *
+ * <p>Package-private access lets concurrency tests submit a deliberately
stale snapshot while
+ * exercising the same root-first ordering as the public delete path.
+ */
+ void deleteViewWithVersion(NameIdentifier identifier, ViewPO observedViewPO)
{
+ SessionUtils.doMultipleWithCommit(
+ // Check the root version before touching relationships. A stale drop
stops here.
+ () ->
+ OccWriteSupport.deleteWithVersion(
+ () ->
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class,
+ mapper ->
+ mapper.softDeleteViewMetasByViewId(
+ observedViewPO.getViewId(),
observedViewPO.getCurrentVersion())),
+ () -> viewWriteFailure(identifier, observedViewPO)),
+ () -> deleteViewDependents(observedViewPO.getViewId()));
}
@Monitored(
@@ -206,56 +246,30 @@ public class ViewMetaService {
return ops;
}
- private boolean deleteView(Long viewId) {
- AtomicInteger deleteResult = new AtomicInteger(0);
- SessionUtils.doMultipleWithCommit(
- () ->
- deleteResult.set(
- SessionUtils.getWithoutCommit(
- ViewMetaMapper.class, mapper ->
mapper.softDeleteViewMetasByViewId(viewId))),
- () -> {
- if (deleteResult.get() > 0) {
- SessionUtils.doWithoutCommit(
- ViewVersionInfoMapper.class,
- mapper -> mapper.softDeleteViewVersionsByViewId(viewId));
- SessionUtils.doWithoutCommit(
- OwnerMetaMapper.class,
- mapper ->
- mapper.softDeleteOwnerRelByMetadataObjectIdAndType(
- viewId, MetadataObject.Type.VIEW.name()));
- SessionUtils.doWithoutCommit(
- SecurableObjectMapper.class,
- mapper ->
- mapper.softDeleteObjectRelsByMetadataObject(
- viewId, MetadataObject.Type.VIEW.name()));
- SessionUtils.doWithoutCommit(
- TagMetadataObjectRelMapper.class,
- mapper ->
- mapper.softDeleteTagMetadataObjectRelsByMetadataObject(
- viewId, MetadataObject.Type.VIEW.name()));
- SessionUtils.doWithoutCommit(
- PolicyMetadataObjectRelMapper.class,
- mapper ->
- mapper.softDeletePolicyMetadataObjectRelsByMetadataObject(
- viewId, MetadataObject.Type.VIEW.name()));
- }
- });
- return deleteResult.get() > 0;
- }
-
+ /**
+ * Builds the next version of a view row.
+ *
+ * <p>The parent IDs are not seeded here: {@link ViewPO#buildViewPO}
resolves them from the new
+ * entity's namespace, which is what carries a cross-schema move.
+ */
private ViewPO updateViewPO(ViewPO oldViewPO, ViewEntity newEntity) {
- Long newVersion = oldViewPO.getLastVersion() + 1;
+ long newVersion = nextVersion(oldViewPO);
ViewPO.ViewPOBuilder builder =
- ViewPO.builder()
- .withMetalakeId(oldViewPO.getMetalakeId())
- .withCatalogId(oldViewPO.getCatalogId())
- .withSchemaId(oldViewPO.getSchemaId())
- .withCurrentVersion(newVersion)
- .withLastVersion(newVersion);
- return buildViewPO(newEntity, builder, newVersion.intValue());
+
ViewPO.builder().withCurrentVersion(newVersion).withLastVersion(newVersion);
+ return buildViewPO(newEntity, builder, (int) newVersion);
}
- private ViewPO getViewPOByIdentifier(NameIdentifier identifier) {
+ /**
+ * Returns the version to write next, one above every version this view has
ever had.
+ *
+ * @param viewPO the view row observed by the caller
+ * @return the next monotonic version
+ */
+ private static long nextVersion(ViewPO viewPO) {
+ return Math.max(viewPO.getCurrentVersion(), viewPO.getLastVersion()) + 1;
+ }
+
+ ViewPO getViewPOByIdentifier(NameIdentifier identifier) {
NameIdentifierUtil.checkView(identifier);
ViewPO viewPO =
SessionUtils.getWithoutCommit(
@@ -276,4 +290,138 @@ public class ViewMetaService {
ViewMetaMapper.class,
mapper -> POStorageReadRouting.listPOs(mapper, namespace, ops,
Entity.EntityType.VIEW));
}
+
+ /**
+ * Writes a new view or replaces the active view selected by natural key or
stable ID. Locking the
+ * root before choosing the next version makes overwrite behavior identical
across databases and
+ * keeps standard reads strict about the matching version-row invariant.
+ */
+ private void insertViewWithoutCommit(
+ ViewEntity viewEntity, ViewPO initializedViewPO, boolean overwrite) {
+ if (!overwrite) {
+ insertNewViewWithoutCommit(initializedViewPO);
+ return;
+ }
+
+ ViewPO existingViewPO = findAndLockViewForOverwrite(initializedViewPO);
+ if (existingViewPO == null) {
+ if (SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class,
+ mapper ->
mapper.countDeletedViewMetasById(initializedViewPO.getViewId()))
+ > 0) {
+ throw new EntityAlreadyExistsException(
+ "The view ID %s is reserved by a deleted view; use a new ID",
+ initializedViewPO.getViewId());
+ }
+ insertNewViewWithoutCommit(initializedViewPO);
+ return;
+ }
+
+ ViewPO replacementPO = viewPOForOverwrite(initializedViewPO,
existingViewPO);
+ int updated =
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class, mapper -> ops.updatePO(mapper,
replacementPO, existingViewPO));
+ if (updated == 0) {
+ throw viewWriteFailure(viewEntity.nameIdentifier(), existingViewPO);
+ }
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.insertViewVersionInfo(replacementPO.getViewVersionInfoPO()));
+ }
+
+ private void insertNewViewWithoutCommit(ViewPO viewPO) {
+ SessionUtils.doWithoutCommit(
+ ViewMetaMapper.class, mapper -> ops.insertPO(mapper, viewPO, false));
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper -> mapper.insertViewVersionInfo(viewPO.getViewVersionInfoPO()));
+ }
+
+ private ViewPO findAndLockViewForOverwrite(ViewPO initializedViewPO) {
+ return OccWriteSupport.findAndLockForOverwrite(
+ () ->
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class,
+ mapper ->
+ mapper.selectViewMetaBySchemaIdAndNameForUpdate(
+ initializedViewPO.getSchemaId(),
initializedViewPO.getViewName())),
+ () ->
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class,
+ mapper ->
mapper.selectViewMetaByIdForUpdate(initializedViewPO.getViewId())),
+ current -> Objects.equals(current.getSchemaId(),
initializedViewPO.getSchemaId()));
+ }
+
+ private ViewPO viewPOForOverwrite(ViewPO incomingPO, ViewPO persistedPO) {
+ Long nextVersion = nextVersion(persistedPO);
+ ViewVersionInfoPO incomingVersionPO = incomingPO.getViewVersionInfoPO();
+ ViewVersionInfoPO persistedVersionPO =
+ ViewVersionInfoPO.builder()
+ .withMetalakeId(incomingVersionPO.metalakeId())
+ .withCatalogId(incomingVersionPO.catalogId())
+ .withSchemaId(incomingVersionPO.schemaId())
+ .withViewId(persistedPO.getViewId())
+ .withVersion(nextVersion.intValue())
+ .withViewComment(incomingVersionPO.viewComment())
+ .withColumns(incomingVersionPO.columns())
+ .withProperties(incomingVersionPO.properties())
+ .withDefaultCatalog(incomingVersionPO.defaultCatalog())
+ .withDefaultSchema(incomingVersionPO.defaultSchema())
+ .withRepresentations(incomingVersionPO.representations())
+ .withAuditInfo(incomingVersionPO.auditInfo())
+ .withDeletedAt(incomingVersionPO.deletedAt())
+ .build();
+ return ViewPO.builder()
+ .withViewId(persistedPO.getViewId())
+ .withViewName(incomingPO.getViewName())
+ .withMetalakeId(incomingPO.getMetalakeId())
+ .withCatalogId(incomingPO.getCatalogId())
+ .withSchemaId(incomingPO.getSchemaId())
+ .withAuditInfo(incomingPO.getAuditInfo())
+ .withCurrentVersion(nextVersion)
+ .withLastVersion(nextVersion)
+ .withDeletedAt(incomingPO.getDeletedAt())
+ .withViewVersionInfoPO(persistedVersionPO)
+ .build();
+ }
+
+ private void deleteViewDependents(Long viewId) {
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class, mapper ->
mapper.softDeleteViewVersionsByViewId(viewId));
+ SessionUtils.doWithoutCommit(
+ OwnerMetaMapper.class,
+ mapper ->
+ mapper.softDeleteOwnerRelByMetadataObjectIdAndType(
+ viewId, MetadataObject.Type.VIEW.name()));
+ SessionUtils.doWithoutCommit(
+ SecurableObjectMapper.class,
+ mapper ->
+ mapper.softDeleteObjectRelsByMetadataObject(viewId,
MetadataObject.Type.VIEW.name()));
+ SessionUtils.doWithoutCommit(
+ TagMetadataObjectRelMapper.class,
+ mapper ->
+ mapper.softDeleteTagMetadataObjectRelsByMetadataObject(
+ viewId, MetadataObject.Type.VIEW.name()));
+ SessionUtils.doWithoutCommit(
+ PolicyMetadataObjectRelMapper.class,
+ mapper ->
+ mapper.softDeletePolicyMetadataObjectRelsByMetadataObject(
+ viewId, MetadataObject.Type.VIEW.name()));
+ }
+
+ private RuntimeException viewWriteFailure(NameIdentifier identifier, ViewPO
observedViewPO) {
+ return OccWriteSupport.writeFailure(
+ identifier,
+ Entity.EntityType.VIEW,
+ () ->
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class,
+ mapper ->
mapper.selectViewMetaByIdForUpdate(observedViewPO.getViewId())),
+ null,
+ current ->
+ Objects.equals(current.getViewName(), observedViewPO.getViewName())
+ && Objects.equals(current.getSchemaId(),
observedViewPO.getSchemaId())
+ && Objects.equals(current.getCatalogId(),
observedViewPO.getCatalogId())
+ && Objects.equals(current.getMetalakeId(),
observedViewPO.getMetalakeId()));
+ }
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewPOStorageOps.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewPOStorageOps.java
index 7294c50887..a7bf91cc56 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewPOStorageOps.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewPOStorageOps.java
@@ -18,6 +18,7 @@
*/
package org.apache.gravitino.storage.relational.service;
+import com.google.common.base.Preconditions;
import java.util.List;
import java.util.stream.Collectors;
import org.apache.gravitino.Entity;
@@ -31,13 +32,19 @@ public class ViewPOStorageOps extends
BasePOStorageOps<ViewPO, ViewMetaMapper> {
public ViewPOStorageOps() {}
+ /**
+ * {@inheritDoc}
+ *
+ * <p>Overwrite is not supported here. {@link ViewMetaService} replaces an
existing view by
+ * locking its root row and advancing it with a version compare-and-set,
which keeps the stored ID
+ * and the version sequence identical on every database. A database-specific
upsert would bypass
+ * that and is rejected instead of silently taking a second code path.
+ */
@Override
public void insertPO(ViewMetaMapper mapper, ViewPO viewPO, boolean
overwrite) {
- if (overwrite) {
- mapper.insertViewMetaOnDuplicateKeyUpdate(viewPO);
- } else {
- mapper.insertViewMeta(viewPO);
- }
+ Preconditions.checkArgument(
+ !overwrite, "View overwrite is handled by ViewMetaService, not by an
upsert");
+ mapper.insertViewMeta(viewPO);
}
@Override
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestViewMetaBaseSQLProvider.java
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestViewMetaBaseSQLProvider.java
new file mode 100644
index 0000000000..36fd8eca23
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestViewMetaBaseSQLProvider.java
@@ -0,0 +1,76 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.gravitino.storage.relational.mapper.provider.base;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+class TestViewMetaBaseSQLProvider {
+
+ private static final ViewMetaBaseSQLProvider PROVIDER = new
ViewMetaBaseSQLProvider();
+
+ @Test
+ void testUpdateUsesVersionCas() {
+ String sql = PROVIDER.updateViewMeta(null, null);
+ Assertions.assertTrue(sql.contains("metalake_id =
#{newViewMeta.metalakeId}"));
+ Assertions.assertTrue(sql.contains("catalog_id =
#{newViewMeta.catalogId}"));
+ String whereClause = sql.substring(sql.indexOf(" WHERE"));
+
+ Assertions.assertEquals(
+ " WHERE view_id = #{oldViewMeta.viewId} "
+ + " AND current_version = #{oldViewMeta.currentVersion} "
+ + " AND deleted_at = 0",
+ whereClause);
+ }
+
+ @Test
+ void testDirectDeleteUsesVersionCas() {
+ String sql = PROVIDER.softDeleteViewMetasByViewId(null, null);
+
+ Assertions.assertTrue(sql.contains("AND current_version =
#{currentVersion}"));
+ Assertions.assertTrue(sql.endsWith("AND deleted_at = 0"));
+ }
+
+ @Test
+ void testConflictReadLocksActiveRow() {
+ String sql = PROVIDER.selectViewMetaByIdForUpdate(null);
+
+ Assertions.assertTrue(sql.contains("view_name as viewName"));
+ Assertions.assertTrue(sql.contains("current_version as currentVersion"));
+ Assertions.assertTrue(sql.contains("WHERE view_id = #{viewId} AND
deleted_at = 0"));
+ Assertions.assertTrue(sql.endsWith("FOR UPDATE"));
+ }
+
+ @Test
+ void testNormalReadRequiresCurrentVersionRow() {
+ String sql = PROVIDER.selectViewMetaBySchemaIdAndName(null, null);
+
+ Assertions.assertTrue(sql.contains("vm INNER JOIN view_version_info vi"));
+ Assertions.assertTrue(sql.endsWith("AND vm.deleted_at = 0 AND
vi.deleted_at = 0"));
+ }
+
+ @Test
+ void testOverwriteLookupLocksOnlyViewRoot() {
+ String sql = PROVIDER.selectViewMetaBySchemaIdAndNameForUpdate(null, null);
+
+ Assertions.assertFalse(sql.contains("view_version_info"));
+ Assertions.assertTrue(sql.contains("current_version as currentVersion"));
+ Assertions.assertTrue(sql.endsWith("AND deleted_at = 0 FOR UPDATE"));
+ }
+}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TestViewMetaPostgreSQLProvider.java
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TestViewMetaPostgreSQLProvider.java
new file mode 100644
index 0000000000..dc036ea419
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TestViewMetaPostgreSQLProvider.java
@@ -0,0 +1,41 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.gravitino.storage.relational.mapper.provider.postgresql;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+class TestViewMetaPostgreSQLProvider {
+
+ @Test
+ void testDirectDeleteUsesVersionCas() {
+ String sql = new
ViewMetaPostgreSQLProvider().softDeleteViewMetasByViewId(null, null);
+
+ Assertions.assertTrue(sql.contains("AND current_version =
#{currentVersion}"));
+ Assertions.assertTrue(sql.endsWith("AND deleted_at = 0"));
+ }
+
+ @Test
+ void testNormalReadRequiresCurrentVersionRow() {
+ String sql = new
ViewMetaPostgreSQLProvider().selectViewMetaBySchemaIdAndName(null, null);
+
+ Assertions.assertTrue(sql.contains("vm INNER JOIN view_version_info vi"));
+ Assertions.assertTrue(sql.endsWith("AND vm.deleted_at = 0 AND
vi.deleted_at = 0"));
+ }
+}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOccWriteSupport.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOccWriteSupport.java
index 8e1ecbf01b..c09c433ad1 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOccWriteSupport.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOccWriteSupport.java
@@ -21,12 +21,15 @@ package org.apache.gravitino.storage.relational.service;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
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.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
import java.util.Collections;
import java.util.List;
import java.util.Objects;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.exceptions.OptimisticLockException;
@@ -52,6 +55,48 @@ public class TestOccWriteSupport {
}
}
+ @Test
+ void testOverwritePrefersNaturalKey() {
+ DummyPO target = new DummyPO("target", 1L);
+ assertSame(
+ target,
+ OccWriteSupport.findAndLockForOverwrite(
+ () -> target,
+ () -> {
+ throw new AssertionError("ID lookup must be skipped");
+ },
+ po -> true));
+ }
+
+ @Test
+ void testOverwriteFallsBackToIdInSameParent() {
+ DummyPO target = new DummyPO("renamed", 1L);
+ assertSame(
+ target,
+ OccWriteSupport.findAndLockForOverwrite(
+ () -> null, () -> target, po -> po.parentId() == 1L));
+ }
+
+ @Test
+ void testOverwriteRejectsIdInDifferentParent() {
+ assertThrows(
+ EntityAlreadyExistsException.class,
+ () ->
+ OccWriteSupport.findAndLockForOverwrite(
+ () -> null, () -> new DummyPO("foreign", 2L), po ->
po.parentId() == 1L));
+ }
+
+ @Test
+ void testOverwriteReturnsNullWhenBothLookupsMiss() {
+ assertNull(
+ OccWriteSupport.findAndLockForOverwrite(
+ () -> null,
+ () -> null,
+ po -> {
+ throw new AssertionError("No row to validate");
+ }));
+ }
+
@Test
void testWriteFailureReturnsNoSuchEntityWhenNotFound() {
NameIdentifier ident = NameIdentifier.of("metalake_test");
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
index b971030a85..74917c154c 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
@@ -93,6 +93,120 @@ public class TestTagMetaService extends TestJDBCBackend {
private final Map<String, String> props = ImmutableMap.of("k1", "v1");
+ /** Verifies a deleted ID is rejected permanently instead of producing a
retryable conflict. */
+ @TestTemplate
+ public void testOverwriteRejectsDeletedTagId() throws IOException {
+ createAndInsertMakeLake(METALAKE_NAME);
+ TagEntity tag =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("deleted_tag_id")
+ .withNamespace(NamespaceUtil.ofTag(METALAKE_NAME))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TagMetaService service = TagMetaService.getInstance();
+ service.insertTag(tag, false);
+ assertTrue(service.deleteTag(tag.nameIdentifier()));
+ EntityAlreadyExistsException failure =
+ Assertions.assertThrows(
+ EntityAlreadyExistsException.class, () -> service.insertTag(tag,
true));
+ assertTrue(failure.getMessage().contains("use a new ID"));
+ Assertions.assertThrows(
+ NoSuchEntityException.class, () ->
service.getTagByIdentifier(tag.nameIdentifier()));
+ TagEntity replacement =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(tag.name())
+ .withNamespace(tag.namespace())
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ service.insertTag(replacement, true);
+ assertEquals(replacement.id(),
service.getTagByIdentifier(tag.nameIdentifier()).id());
+ }
+
+ /** Verifies first-time overwrites either serialize or report a retryable
insert conflict. */
+ @TestTemplate
+ public void testConcurrentOverwriteOfMissingTag() throws Exception {
+ createAndInsertMakeLake(METALAKE_NAME);
+ TagEntity first =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("tag_first_overwrite")
+ .withNamespace(NamespaceUtil.ofTag(METALAKE_NAME))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TagEntity second =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(first.name())
+ .withNamespace(first.namespace())
+ .withComment("second overwrite")
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TagMetaService service = TagMetaService.getInstance();
+ CountDownLatch firstWritten = new CountDownLatch(1);
+ CountDownLatch allowCommit = new CountDownLatch(1);
+ CountDownLatch secondStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> firstResult =
+ executor.submit(
+ () -> {
+ SessionUtils.beginTransaction();
+ try {
+ service.insertTag(first, true);
+ firstWritten.countDown();
+ await(allowCommit);
+ SessionUtils.commitTransaction();
+ return null;
+ } catch (Throwable failure) {
+ SessionUtils.rollbackTransaction();
+ return failure;
+ }
+ });
+ try {
+ assertTrue(firstWritten.await(30, TimeUnit.SECONDS));
+ Future<Throwable> secondResult =
+ executor.submit(
+ () -> {
+ secondStarted.countDown();
+ try {
+ service.insertTag(second, true);
+ return null;
+ } catch (Throwable failure) {
+ return failure;
+ }
+ });
+ assertTrue(secondStarted.await(30, TimeUnit.SECONDS));
+ Assertions.assertThrows(
+ TimeoutException.class, () -> secondResult.get(500,
TimeUnit.MILLISECONDS));
+ allowCommit.countDown();
+ Assertions.assertNull(firstResult.get(30, TimeUnit.SECONDS));
+ Throwable failure = secondResult.get(30, TimeUnit.SECONDS);
+ if (failure != null) {
+ assertTrue(
+ failure instanceof OptimisticLockException, () -> "Unexpected
failure: " + failure);
+ assertEquals(first.id(),
service.getTagByIdentifier(first.nameIdentifier()).id());
+ Assertions.assertNull(
+ SessionUtils.getWithoutCommit(
+ TagMetaMapper.class, mapper ->
mapper.selectTagByTagId(second.id())));
+ service.insertTag(second, true);
+ }
+ TagEntity stored = service.getTagByIdentifier(first.nameIdentifier());
+ assertEquals(first.id(), stored.id());
+ assertEquals("second overwrite", stored.comment());
+ TagPO storedPO =
+ SessionUtils.getWithoutCommit(
+ TagMetaMapper.class, mapper ->
mapper.selectTagByTagId(first.id()));
+ assertEquals(2L, storedPO.getCurrentVersion().longValue());
+ Assertions.assertNull(
+ SessionUtils.getWithoutCommit(
+ TagMetaMapper.class, mapper ->
mapper.selectTagByTagId(second.id())));
+ } finally {
+ allowCommit.countDown();
+ executor.shutdownNow();
+ }
+ }
+
@TestTemplate
public void testMetaLifeCycleFromCreationToDeletion() throws IOException {
BaseMetalake metalake = createAndInsertMakeLake(METALAKE_NAME);
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestViewMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestViewMetaService.java
index edeaff2565..9e2500ad82 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestViewMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestViewMetaService.java
@@ -20,6 +20,7 @@ package org.apache.gravitino.storage.relational.service;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -33,10 +34,17 @@ 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.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.AuditInfo;
import org.apache.gravitino.meta.TagEntity;
@@ -47,7 +55,13 @@ import org.apache.gravitino.rel.SQLRepresentation;
import org.apache.gravitino.rel.types.Types;
import org.apache.gravitino.storage.RandomIdGenerator;
import org.apache.gravitino.storage.relational.TestJDBCBackend;
+import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ViewVersionInfoMapper;
+import org.apache.gravitino.storage.relational.po.SchemaPO;
+import org.apache.gravitino.storage.relational.po.ViewPO;
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;
@@ -67,6 +81,22 @@ public class TestViewMetaService extends TestJDBCBackend {
createAndInsertSchema(metalakeName, catalogName, schemaName);
}
+ /** Dropping the old parent after a move must preserve every historical
version. */
+ @TestTemplate
+ public void testMovedViewHistorySurvivesSourceCascade() throws IOException {
+ for (String parent : new String[] {"schema", "catalog", "metalake"}) {
+ assertMovedViewHistoryCascade(parent, true);
+ }
+ }
+
+ /** Dropping the current parent must delete versions created under previous
parents too. */
+ @TestTemplate
+ public void testMovedViewHistoryIsDeletedWithDestination() throws
IOException {
+ for (String parent : new String[] {"schema", "catalog", "metalake"}) {
+ assertMovedViewHistoryCascade(parent, false);
+ }
+ }
+
@TestTemplate
public void testInsertAlreadyExistsException() throws IOException {
Namespace ns = NamespaceUtil.ofView(metalakeName, catalogName, schemaName);
@@ -162,7 +192,80 @@ public class TestViewMetaService extends TestJDBCBackend {
}
@TestTemplate
- public void testUpdateViewRollbackWhenMetaUpdateAffectsZeroRows() throws
IOException {
+ public void testCreateViewWaitsForConcurrentSchemaDelete() throws Exception {
+ SchemaPO observedSchemaPO =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+ mapper.selectSchemaByFullQualifiedName(metalakeName,
catalogName, schemaName));
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity view =
+ createViewEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ GravitinoITUtils.genRandomName("view_parent_delete_race"),
+ AUDIT_INFO);
+
+ CountDownLatch deleteWritten = new CountDownLatch(1);
+ CountDownLatch allowDeleteCommit = new CountDownLatch(1);
+ CountDownLatch createStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> deleteResult =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ assertEquals(
+ Integer.valueOf(1),
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+
mapper.softDeleteSchemaMetaBySchemaIdAndVersion(
+ observedSchemaPO.getSchemaId(),
+
observedSchemaPO.getCurrentVersion()))),
+ () -> {
+ deleteWritten.countDown();
+ await(allowDeleteCommit);
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ try {
+ assertTrue(deleteWritten.await(30, TimeUnit.SECONDS));
+ Future<Throwable> createResult =
+ executor.submit(
+ () -> {
+ createStarted.countDown();
+ try {
+ ViewMetaService.getInstance().insertView(view, false);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ assertTrue(createStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> createResult.get(500,
TimeUnit.MILLISECONDS));
+
+ allowDeleteCommit.countDown();
+ assertNull(deleteResult.get(30, TimeUnit.SECONDS));
+ Throwable createFailure = createResult.get(30, TimeUnit.SECONDS);
+ assertTrue(createFailure instanceof NoSuchEntityException,
String.valueOf(createFailure));
+ } finally {
+ allowDeleteCommit.countDown();
+ executor.shutdownNow();
+ }
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier()));
+ }
+
+ @TestTemplate
+ public void testUpdateViewReportsNoSuchAfterConcurrentDelete() throws
IOException {
String viewName = GravitinoITUtils.genRandomName("test_view");
Namespace ns = NamespaceUtil.ofView(metalakeName, catalogName, schemaName);
ViewEntity view =
@@ -185,7 +288,7 @@ public class TestViewMetaService extends TestJDBCBackend {
.build();
assertThrows(
- IOException.class,
+ NoSuchEntityException.class,
() ->
ViewMetaService.getInstance()
.updateView(
@@ -201,6 +304,292 @@ public class TestViewMetaService extends TestJDBCBackend {
assertTrue(versions.get(1) > 0L);
}
+ @TestTemplate
+ public void testAlterReportsOptimisticLockConflictAndKeepsWinnerVersion()
throws IOException {
+ String viewName = GravitinoITUtils.genRandomName("view_alter_conflict");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+
+ assertThrows(
+ OptimisticLockException.class,
+ () ->
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ entity -> {
+ try {
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ competing ->
+ copyViewWithComment(
+ (ViewEntity) competing, "competing
update"));
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ return copyViewWithComment((ViewEntity) entity,
"requested update");
+ }));
+
+ ViewEntity current =
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier());
+ assertEquals("competing update", current.comment());
+ Map<Integer, Long> versions = listViewVersions(view.id());
+ assertEquals(2, versions.size());
+ assertTrue(versions.containsKey(2));
+ }
+
+ @TestTemplate
+ public void testAlterRollsBackRootCasWhenVersionInsertFails() throws
IOException {
+ String viewName =
GravitinoITUtils.genRandomName("view_version_insert_failure");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+ ViewEntity conflictingVersion = copyViewWithComment(view, "conflicting
version");
+ ViewPO conflictingPO =
+ ViewPO.buildViewPO(
+ conflictingVersion,
ViewPO.builder().withCurrentVersion(2L).withLastVersion(2L), 2);
+ SessionUtils.doWithCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.insertViewVersionInfo(conflictingPO.getViewVersionInfoPO()));
+
+ assertThrows(
+ EntityAlreadyExistsException.class,
+ () ->
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ entity -> copyViewWithComment((ViewEntity) entity, "must
roll back")));
+
+ ViewPO currentPO =
ViewMetaService.getInstance().getViewPOByIdentifier(view.nameIdentifier());
+ assertEquals(1L, currentPO.getCurrentVersion());
+ assertEquals(1L, currentPO.getLastVersion());
+ assertEquals(
+ view.comment(),
+
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier()).comment());
+ }
+
+ @TestTemplate
+ public void testAlterReportsNoSuchWhenRenamedConcurrently() throws
IOException {
+ String viewName = GravitinoITUtils.genRandomName("view_alter_renamed");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+ NameIdentifier renamedIdentifier = NameIdentifier.of(namespace, viewName +
"_winner");
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ entity -> {
+ try {
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ competing ->
+ copyView(
+ (ViewEntity) competing,
+ renamedIdentifier.name(),
+ namespace,
+ "renamed winner"));
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ return copyViewWithComment((ViewEntity) entity, "stale
update");
+ }));
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier()));
+ assertEquals(
+ "renamed winner",
+
ViewMetaService.getInstance().getViewByIdentifier(renamedIdentifier).comment());
+ }
+
+ @TestTemplate
+ public void testAlterReportsNoSuchWhenMovedConcurrently() throws IOException
{
+ String targetCatalogName =
GravitinoITUtils.genRandomName("view_target_catalog");
+ String targetSchemaName =
GravitinoITUtils.genRandomName("view_target_schema");
+ createAndInsertCatalog(metalakeName, targetCatalogName);
+ createAndInsertSchema(metalakeName, targetCatalogName, targetSchemaName);
+ String viewName = GravitinoITUtils.genRandomName("view_alter_moved");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ Namespace movedNamespace =
+ NamespaceUtil.ofView(metalakeName, targetCatalogName,
targetSchemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+ NameIdentifier movedIdentifier = NameIdentifier.of(movedNamespace,
viewName);
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ entity -> {
+ try {
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ competing ->
+ copyView(
+ (ViewEntity) competing,
+ viewName,
+ movedNamespace,
+ "moved winner"));
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ return copyViewWithComment((ViewEntity) entity, "stale
update");
+ }));
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier()));
+ assertEquals(
+ "moved winner",
+
ViewMetaService.getInstance().getViewByIdentifier(movedIdentifier).comment());
+ }
+
+ @TestTemplate
+ public void testMoveWaitsForConcurrentTargetSchemaDelete() throws Exception {
+ String targetCatalogName =
GravitinoITUtils.genRandomName("view_deleted_target_catalog");
+ String targetSchemaName =
GravitinoITUtils.genRandomName("view_deleted_target_schema");
+ createAndInsertCatalog(metalakeName, targetCatalogName);
+ createAndInsertSchema(metalakeName, targetCatalogName, targetSchemaName);
+ SchemaPO targetSchemaPO =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+ mapper.selectSchemaByFullQualifiedName(
+ metalakeName, targetCatalogName, targetSchemaName));
+ String viewName =
GravitinoITUtils.genRandomName("view_target_delete_race");
+ Namespace sourceNamespace = NamespaceUtil.ofView(metalakeName,
catalogName, schemaName);
+ Namespace targetNamespace =
+ NamespaceUtil.ofView(metalakeName, targetCatalogName,
targetSchemaName);
+ ViewEntity view =
+ createViewEntity(
+ RandomIdGenerator.INSTANCE.nextId(), sourceNamespace, viewName,
AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+ ViewEntity moved = copyView(view, viewName, targetNamespace, "must not
move");
+
+ CountDownLatch deleteWritten = new CountDownLatch(1);
+ CountDownLatch allowDeleteCommit = new CountDownLatch(1);
+ CountDownLatch moveStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> deleteResult =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ assertEquals(
+ Integer.valueOf(1),
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+
mapper.softDeleteSchemaMetaBySchemaIdAndVersion(
+ targetSchemaPO.getSchemaId(),
+ targetSchemaPO.getCurrentVersion()))),
+ () -> {
+ deleteWritten.countDown();
+ await(allowDeleteCommit);
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ try {
+ assertTrue(deleteWritten.await(30, TimeUnit.SECONDS));
+ Future<Throwable> moveResult =
+ executor.submit(
+ () -> {
+ moveStarted.countDown();
+ try {
+
ViewMetaService.getInstance().updateView(view.nameIdentifier(), ignored ->
moved);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ assertTrue(moveStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> moveResult.get(500,
TimeUnit.MILLISECONDS));
+
+ allowDeleteCommit.countDown();
+ assertNull(deleteResult.get(30, TimeUnit.SECONDS));
+ Throwable moveFailure = moveResult.get(30, TimeUnit.SECONDS);
+ assertTrue(moveFailure instanceof NoSuchEntityException,
String.valueOf(moveFailure));
+ } finally {
+ allowDeleteCommit.countDown();
+ executor.shutdownNow();
+ }
+
+ assertEquals(
+ view.comment(),
+
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier()).comment());
+ assertEquals(1, listViewVersions(view.id()).size());
+ }
+
+ @TestTemplate
+ public void testDeleteRejectsStaleVersion() throws IOException {
+ String viewName = GravitinoITUtils.genRandomName("view_stale_delete");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+ TagEntity tag =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("view_occ_tag")
+ .withNamespace(NamespaceUtil.ofTag(metalakeName))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TagMetaService.getInstance().insertTag(tag, false);
+ TagMetaService.getInstance()
+ .associateTagsWithMetadataObject(
+ view.nameIdentifier(),
+ view.type(),
+ new NameIdentifier[] {tag.nameIdentifier()},
+ new NameIdentifier[0]);
+ ViewPO stalePO =
ViewMetaService.getInstance().getViewPOByIdentifier(view.nameIdentifier());
+
+ ViewMetaService.getInstance()
+ .updateView(
+ view.nameIdentifier(),
+ entity -> copyViewWithComment((ViewEntity) entity, "winning
update"));
+
+ assertThrows(
+ OptimisticLockException.class,
+ () ->
ViewMetaService.getInstance().deleteViewWithVersion(view.nameIdentifier(),
stalePO));
+ assertEquals(
+ "winning update",
+
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier()).comment());
+ assertEquals(1, countActiveTagRelForMetadataObject(view.id(), "VIEW"));
+ }
+
+ @TestTemplate
+ public void testDeleteReportsNoSuchWhenDeletedConcurrently() throws
IOException {
+ String viewName = GravitinoITUtils.genRandomName("view_delete_deleted");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+ ViewPO stalePO =
ViewMetaService.getInstance().getViewPOByIdentifier(view.nameIdentifier());
+
+ ViewMetaService.getInstance().deleteView(view.nameIdentifier());
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
ViewMetaService.getInstance().deleteViewWithVersion(view.nameIdentifier(),
stalePO));
+ }
+
@TestTemplate
public void testDeleteView() throws IOException {
String viewName = GravitinoITUtils.genRandomName("test_view");
@@ -273,6 +662,325 @@ public class TestViewMetaService extends TestJDBCBackend {
NameIdentifier viewIdent = NameIdentifier.of(metalakeName, catalogName,
schemaName, viewName);
ViewEntity loaded =
ViewMetaService.getInstance().getViewByIdentifier(viewIdent);
assertEquals("overwritten comment", loaded.comment());
+ ViewPO overwrittenPO =
ViewMetaService.getInstance().getViewPOByIdentifier(viewIdent);
+ assertEquals(2L, overwrittenPO.getCurrentVersion());
+ assertEquals(2L, overwrittenPO.getLastVersion());
+ assertEquals(2, listViewVersions(view.id()).size());
+ }
+
+ /** Moving a view must persist every parent ID used by catalog cascade
cleanup. */
+ @TestTemplate
+ public void testMoveAcrossMetalakesPreservesOwnership() throws IOException {
+ String targetMetalake = "moved_view_metalake";
+ String targetCatalog = "moved_view_catalog";
+ String targetSchema = "moved_view_schema";
+ long targetMetalakeId = createAndInsertMakeLake(targetMetalake).id();
+ long targetCatalogId = createAndInsertCatalog(targetMetalake,
targetCatalog).id();
+ createAndInsertSchema(targetMetalake, targetCatalog, targetSchema);
+
+ ViewMetaService service = ViewMetaService.getInstance();
+ ViewEntity original =
+ createViewEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofView(metalakeName, catalogName, schemaName),
+ "moved_view",
+ AUDIT_INFO);
+ service.insertView(original, false);
+ ViewEntity moved =
+ copyView(
+ original,
+ original.name(),
+ NamespaceUtil.ofView(targetMetalake, targetCatalog, targetSchema),
+ "moved");
+ service.updateView(original.nameIdentifier(), ignored -> moved);
+ ViewPO persisted = service.getViewPOByIdentifier(moved.nameIdentifier());
+ assertEquals(targetMetalakeId, persisted.getMetalakeId());
+ assertEquals(targetCatalogId, persisted.getCatalogId());
+ assertEquals(targetMetalakeId,
persisted.getViewVersionInfoPO().metalakeId());
+ assertEquals(targetCatalogId,
persisted.getViewVersionInfoPO().catalogId());
+ assertThrows(
+ NoSuchEntityException.class, () ->
service.getViewByIdentifier(original.nameIdentifier()));
+
+ assertTrue(
+ CatalogMetaService.getInstance()
+ .deleteCatalog(NameIdentifier.of(metalakeName, catalogName),
true));
+ assertEquals(original.id(),
service.getViewByIdentifier(moved.nameIdentifier()).id());
+ assertTrue(
+ CatalogMetaService.getInstance()
+ .deleteCatalog(NameIdentifier.of(targetMetalake, targetCatalog),
true));
+ assertThrows(
+ NoSuchEntityException.class, () ->
service.getViewByIdentifier(moved.nameIdentifier()));
+ listViewVersions(original.id()).values().forEach(deletedAt ->
assertTrue(deletedAt > 0));
+ }
+
+ @TestTemplate
+ public void testCascadingSchemaDeleteCleansViewVersions() throws IOException
{
+ String cascadeSchemaName =
GravitinoITUtils.genRandomName("tst_view_schema_cascade");
+ createAndInsertSchema(metalakeName, catalogName, cascadeSchemaName);
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
cascadeSchemaName);
+ ViewEntity view =
+ createViewEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ GravitinoITUtils.genRandomName("view_cascade"),
+ AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+ ViewMetaService.getInstance()
+ .updateView(view.nameIdentifier(), e ->
copyViewWithComment((ViewEntity) e, "v2"));
+ assertEquals(2, listViewVersions(view.id()).size());
+
+ assertTrue(
+ SchemaMetaService.getInstance()
+ .deleteSchema(NameIdentifier.of(metalakeName, catalogName,
cascadeSchemaName), true));
+
+ // The view root and its version rows must go together. Leaving the
versions active would leak
+ // them: the legacy-timeline collector only reclaims rows that are already
soft deleted.
+ listViewVersions(view.id())
+ .forEach((version, deletedAt) -> assertTrue(deletedAt > 0L, "version "
+ version));
+ }
+
+ @TestTemplate
+ public void testOverwriteDoesNotAdoptViewFromAnotherSchema() throws
IOException {
+ String otherSchemaName =
GravitinoITUtils.genRandomName("tst_view_schema_other");
+ createAndInsertSchema(metalakeName, catalogName, otherSchemaName);
+
+ String viewName =
GravitinoITUtils.genRandomName("view_cross_schema_overwrite");
+ Namespace otherNamespace = NamespaceUtil.ofView(metalakeName, catalogName,
otherSchemaName);
+ ViewEntity foreign =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), otherNamespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(foreign, false);
+
+ // The overwrite targets a schema that holds no view with this name, but
it carries the ID of a
+ // view stored under another schema. Adopting that row would move it out
of its own schema
+ // without ever fencing that schema, so the write is rejected instead.
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity replacement = copyView(foreign, viewName, namespace,
"replacement");
+ assertThrows(
+ EntityAlreadyExistsException.class,
+ () -> ViewMetaService.getInstance().insertView(replacement, true));
+
+ ViewEntity stored =
ViewMetaService.getInstance().getViewByIdentifier(foreign.nameIdentifier());
+ ViewPO storedPO =
ViewMetaService.getInstance().getViewPOByIdentifier(foreign.nameIdentifier());
+ assertEquals(otherNamespace, stored.namespace());
+ assertEquals(foreign.comment(), stored.comment());
+ assertEquals(1L, storedPO.getCurrentVersion());
+ assertEquals(1, listViewVersions(foreign.id()).size());
+ }
+
+ @TestTemplate
+ public void testNaturalKeyOverwriteUsesPersistedViewId() throws IOException {
+ String viewName =
GravitinoITUtils.genRandomName("view_natural_key_overwrite");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity original =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(original, false);
+ ViewEntity replacement =
+ copyView(
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO),
+ viewName,
+ namespace,
+ "replacement");
+
+ ViewMetaService.getInstance().insertView(replacement, true);
+
+ ViewEntity stored =
+
ViewMetaService.getInstance().getViewByIdentifier(original.nameIdentifier());
+ ViewPO storedPO =
+
ViewMetaService.getInstance().getViewPOByIdentifier(original.nameIdentifier());
+ assertEquals(original.id(), stored.id());
+ assertEquals("replacement", stored.comment());
+ assertEquals(2L, storedPO.getCurrentVersion());
+ assertEquals(2, listViewVersions(original.id()).size());
+ assertTrue(listViewVersions(replacement.id()).isEmpty());
+ }
+
+ @TestTemplate
+ public void testNormalReadRequiresCurrentVersionRow() throws IOException {
+ String viewName =
GravitinoITUtils.genRandomName("view_missing_current_version");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(view, false);
+
+ SessionUtils.doWithCommit(
+ ViewVersionInfoMapper.class, mapper ->
mapper.softDeleteViewVersionsByViewId(view.id()));
+
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
ViewMetaService.getInstance().getViewByIdentifier(view.nameIdentifier()));
+ }
+
+ /** A deleted ID is a permanent conflict until GC, rather than a retryable
concurrent write. */
+ @TestTemplate
+ public void testOverwriteRejectsDeletedViewId() throws IOException {
+ Namespace ns = NamespaceUtil.ofView(metalakeName, catalogName, schemaName);
+ ViewEntity view =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), ns,
"deleted_view_id", AUDIT_INFO);
+ ViewMetaService service = ViewMetaService.getInstance();
+ service.insertView(view, false);
+ assertTrue(service.deleteView(view.nameIdentifier()));
+
+ EntityAlreadyExistsException failure =
+ assertThrows(EntityAlreadyExistsException.class, () ->
service.insertView(view, true));
+ assertTrue(failure.getMessage().contains("use a new ID"));
+ assertThrows(
+ NoSuchEntityException.class, () ->
service.getViewByIdentifier(view.nameIdentifier()));
+ listViewVersions(view.id()).values().forEach(deletedAt ->
assertTrue(deletedAt > 0));
+
+ ViewEntity replacement =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), ns, view.name(),
AUDIT_INFO);
+ service.insertView(replacement, true);
+ assertEquals(replacement.id(),
service.getViewByIdentifier(view.nameIdentifier()).id());
+ }
+
+ /** Verifies first-time overwrites either serialize or report a retryable
insert conflict. */
+ @TestTemplate
+ public void testConcurrentOverwriteOfMissingView() throws Exception {
+ String name = GravitinoITUtils.genRandomName("view_first_overwrite");
+ Namespace ns = NamespaceUtil.ofView(metalakeName, catalogName, schemaName);
+ ViewEntity first = createViewEntity(RandomIdGenerator.INSTANCE.nextId(),
ns, name, AUDIT_INFO);
+ ViewEntity second =
+ copyView(
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), ns, name,
AUDIT_INFO),
+ name,
+ ns,
+ "second overwrite");
+ CountDownLatch firstWritten = new CountDownLatch(1);
+ CountDownLatch allowCommit = new CountDownLatch(1);
+ CountDownLatch secondStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> firstResult =
+ executor.submit(
+ () -> {
+ SessionUtils.beginTransaction();
+ try {
+ ViewMetaService.getInstance().insertView(first, true);
+ firstWritten.countDown();
+ await(allowCommit);
+ SessionUtils.commitTransaction();
+ return null;
+ } catch (Throwable failure) {
+ SessionUtils.rollbackTransaction();
+ return failure;
+ }
+ });
+ try {
+ assertTrue(firstWritten.await(30, TimeUnit.SECONDS));
+ Future<Throwable> secondResult =
+ executor.submit(
+ () -> {
+ secondStarted.countDown();
+ try {
+ ViewMetaService.getInstance().insertView(second, true);
+ return null;
+ } catch (Throwable failure) {
+ return failure;
+ }
+ });
+ assertTrue(secondStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> secondResult.get(500,
TimeUnit.MILLISECONDS));
+ allowCommit.countDown();
+ assertNull(firstResult.get(30, TimeUnit.SECONDS));
+ Throwable failure = secondResult.get(30, TimeUnit.SECONDS);
+ // H2 may serialize at the parent lock; MySQL/PostgreSQL can reach the
missing-row INSERT.
+ if (failure != null) {
+ assertTrue(
+ failure instanceof OptimisticLockException, () -> "Unexpected
failure: " + failure);
+ assertEquals(1, listViewVersions(first.id()).size());
+ assertTrue(listViewVersions(second.id()).isEmpty());
+ ViewMetaService.getInstance().insertView(second, true);
+ }
+ ViewEntity stored =
ViewMetaService.getInstance().getViewByIdentifier(first.nameIdentifier());
+ assertEquals(first.id(), stored.id());
+ assertEquals("second overwrite", stored.comment());
+ assertEquals(
+ 2,
+ ViewMetaService.getInstance()
+ .getViewPOByIdentifier(first.nameIdentifier())
+ .getCurrentVersion());
+ assertEquals(2, listViewVersions(first.id()).size());
+ assertTrue(listViewVersions(second.id()).isEmpty());
+ } finally {
+ allowCommit.countDown();
+ executor.shutdownNow();
+ }
+ }
+
+ @TestTemplate
+ public void testNaturalKeyOverwriteWaitsForConcurrentRename() throws
Exception {
+ String viewName =
GravitinoITUtils.genRandomName("view_overwrite_rename_race");
+ Namespace namespace = NamespaceUtil.ofView(metalakeName, catalogName,
schemaName);
+ ViewEntity original =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+ ViewMetaService.getInstance().insertView(original, false);
+ ViewPO observedPO =
+
ViewMetaService.getInstance().getViewPOByIdentifier(original.nameIdentifier());
+ ViewEntity renamed = copyView(original, viewName + "_winner", namespace,
"rename winner");
+ ViewPO renamedPO =
+ ViewPO.buildViewPO(renamed,
ViewPO.builder().withCurrentVersion(2L).withLastVersion(2L), 2);
+ ViewEntity replacement =
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace,
viewName, AUDIT_INFO);
+
+ CountDownLatch renameWritten = new CountDownLatch(1);
+ CountDownLatch allowRenameCommit = new CountDownLatch(1);
+ CountDownLatch overwriteStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> renameResult =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ assertEquals(
+ Integer.valueOf(1),
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class,
+ mapper -> mapper.updateViewMeta(renamedPO,
observedPO))),
+ () ->
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
+
mapper.insertViewVersionInfo(renamedPO.getViewVersionInfoPO())),
+ () -> {
+ renameWritten.countDown();
+ await(allowRenameCommit);
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ try {
+ assertTrue(renameWritten.await(30, TimeUnit.SECONDS));
+ Future<Throwable> overwriteResult =
+ executor.submit(
+ () -> {
+ overwriteStarted.countDown();
+ try {
+ ViewMetaService.getInstance().insertView(replacement, true);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ assertTrue(overwriteStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> overwriteResult.get(500,
TimeUnit.MILLISECONDS));
+
+ allowRenameCommit.countDown();
+ assertNull(renameResult.get(30, TimeUnit.SECONDS));
+ assertNull(overwriteResult.get(30, TimeUnit.SECONDS));
+ } finally {
+ allowRenameCommit.countDown();
+ executor.shutdownNow();
+ }
+
+ assertEquals(
+ original.id(),
+
ViewMetaService.getInstance().getViewByIdentifier(renamed.nameIdentifier()).id());
+ assertEquals(
+ replacement.id(),
+
ViewMetaService.getInstance().getViewByIdentifier(replacement.nameIdentifier()).id());
}
@TestTemplate
@@ -334,6 +1042,99 @@ public class TestViewMetaService extends TestJDBCBackend {
.build();
}
+ private ViewEntity copyViewWithComment(ViewEntity view, String comment) {
+ return copyView(view, view.name(), view.namespace(), comment);
+ }
+
+ private ViewEntity copyView(ViewEntity view, String name, Namespace
namespace, String comment) {
+ return ViewEntity.builder()
+ .withId(view.id())
+ .withName(name)
+ .withNamespace(namespace)
+ .withComment(comment)
+ .withColumns(view.columns())
+ .withRepresentations(view.representations())
+ .withDefaultCatalog(view.defaultCatalog())
+ .withDefaultSchema(view.defaultSchema())
+ .withProperties(view.properties())
+ .withAuditInfo(view.auditInfo())
+ .build();
+ }
+
+ private void await(CountDownLatch latch) {
+ try {
+ assertTrue(latch.await(30, TimeUnit.SECONDS));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ }
+
+ private void assertMovedViewHistoryCascade(String parent, boolean
sourceFirst)
+ throws IOException {
+ String sourceMetalake = "source_view_" + parent + "_" + sourceFirst;
+ String destinationMetalake = "destination_view_" + parent + "_" +
sourceFirst;
+ String catalog = "history_catalog";
+ String schema = "history_schema";
+ for (String metalake : new String[] {sourceMetalake, destinationMetalake})
{
+ createAndInsertMakeLake(metalake);
+ createAndInsertCatalog(metalake, catalog);
+ createAndInsertSchema(metalake, catalog, schema);
+ }
+ ViewMetaService service = ViewMetaService.getInstance();
+ ViewEntity original =
+ createViewEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofView(sourceMetalake, catalog, schema),
+ "moved_history",
+ AUDIT_INFO);
+ service.insertView(original, false);
+ ViewEntity moved =
+ copyView(
+ original,
+ original.name(),
+ NamespaceUtil.ofView(destinationMetalake, catalog, schema),
+ "moved");
+ service.updateView(original.nameIdentifier(), ignored -> moved);
+ assertEquals(2, listViewVersions(original.id()).size());
+
+ if (sourceFirst) {
+ deleteHistoryParent(parent, sourceMetalake, catalog, schema);
+ assertEquals(original.id(),
service.getViewByIdentifier(moved.nameIdentifier()).id());
+ assertEquals(2, listViewVersions(original.id()).size());
+ listViewVersions(original.id())
+ .values()
+ .forEach(deletedAt -> assertEquals(0L, deletedAt.longValue()));
+ }
+
+ deleteHistoryParent(parent, destinationMetalake, catalog, schema);
+ assertThrows(
+ NoSuchEntityException.class, () ->
service.getViewByIdentifier(moved.nameIdentifier()));
+ assertEquals(2, listViewVersions(original.id()).size());
+ listViewVersions(original.id()).values().forEach(deletedAt ->
assertTrue(deletedAt > 0));
+ }
+
+ private void deleteHistoryParent(String parent, String metalake, String
catalog, String schema) {
+ switch (parent) {
+ case "schema":
+ assertTrue(
+ SchemaMetaService.getInstance()
+ .deleteSchema(NameIdentifier.of(metalake, catalog, schema),
true));
+ break;
+ case "catalog":
+ assertTrue(
+ CatalogMetaService.getInstance()
+ .deleteCatalog(NameIdentifier.of(metalake, catalog), true));
+ break;
+ case "metalake":
+ assertTrue(
+
MetalakeMetaService.getInstance().deleteMetalake(NameIdentifier.of(metalake),
true));
+ break;
+ default:
+ throw new AssertionError("Unexpected parent: " + parent);
+ }
+ }
+
private Map<Integer, Long> listViewVersions(Long viewId) {
Map<Integer, Long> versionDeletedTime = new HashMap<>();
try (SqlSession sqlSession =