This is an automated email from the ASF dual-hosted git repository.
yuqi1129 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 3e49a62bb7 [#12406] fix(core): acquire schema lock in child update
paths and clean up table_version_info in schema cascade (#13094)
3e49a62bb7 is described below
commit 3e49a62bb7dce3b7d8d1ac11245018927e9db7f0
Author: li jie <[email protected]>
AuthorDate: Tue Sep 22 17:10:04 2026 +0800
[#12406] fix(core): acquire schema lock in child update paths and clean up
table_version_info in schema cascade (#13094)
### What changes were proposed in this pull request?
Schema cascade delete can race with concurrent child metadata writes and
leave active orphan rows in version tables. This change ensures every
child update transaction acquires the parent schema shared lock
(`lockSchemaForEntityWrite`) as its first operation, so a cascade delete
cannot interleave with an in-flight child write.
Changes:
- **TableMetaService.updateTable**: Always lock the parent schema row
unconditionally (previously only locked when `isSchemaChanged=true`).
For cross-schema moves, lock both source and destination schemas in a
consistent order (by schemaId) to prevent deadlocks.
- **FunctionMetaService.updateFunction**: Same pattern — unconditional
schema lock for same-schema updates, dual-lock for cross-schema moves.
- **ViewMetaService.updateView**: Same pattern.
- **FilesetMetaService.updateFileset**: Add schema lock (previously had
none). Inline the old `tryUpdateFileset` helper into
`doMultipleWithCommit` with lock as the first step, and deduplicate
`retryFilesetPO` construction (was built twice; now built once and
shared by both lambdas).
- **TopicMetaService.updateTopic**: Add schema lock (previously had
none).
- **ModelMetaService.updateModel**: Add schema lock (previously had
none; only `insertModel` had it).
New tests in `TestSchemaMetaService`:
1. `testSchemaChildUpdateServicesWaitForConcurrentSchemaDelete` — runs
cascade delete concurrently with each child type's update (table, view,
fileset, function, model, topic) and asserts the schema and child entity
are gone afterward.
2. `testCascadeDeleteLeavesNoOrphanVersionRows` — runs cascade delete
concurrently with table update, asserts no active `table_version_info`
rows whose parent `table_meta` is deleted (skipped on H2 which lacks
`FOR SHARE`; enforced on MySQL/PostgreSQL).
### Why are the changes needed?
The insert paths (`insertTable`, `insertFunction`, `insertView`, etc.)
already acquire a shared lock on the parent schema row before writing.
However, the corresponding update paths did not consistently do so:
- `updateTable`/`updateFunction`/`updateView`: only locked the schema
when the entity moved to a different schema (`isSchemaChanged=true`);
same-schema updates skipped the lock entirely.
- `updateFileset`/`updateTopic`/`updateModel`: had no schema lock at
all.
As a result, a schema cascade delete could commit while an overlapping
writer left an active row in a child or version table whose schema or
parent entity had already been deleted (orphan row).
Fix: #12406
### Does this PR introduce _any_ user-facing change?
No. The fix is an internal concurrency-safety improvement to the
relational storage layer. No public API, REST endpoint, or behavior
contract changes.
### How was this patch tested?
- `./gradlew :core:spotlessApply :core:compileJava :core:compileTestJava
-PskipWeb=true`
- `./gradlew :core:test --tests
"org.apache.gravitino.storage.relational.service.TestSchemaMetaService"
--tests
"org.apache.gravitino.storage.relational.service.TestModelMetaService"
--tests
"org.apache.gravitino.storage.relational.service.TestFilesetMetaService"
-PskipWeb=true -PskipDockerTests=true`
All existing and new tests pass on H2. The no-orphan-row assertion in
`testCascadeDeleteLeavesNoOrphanVersionRows` is gated to non-H2 backends
because H2 lacks `FOR SHARE` (falls back to `FOR UPDATE`); MySQL and
PostgreSQL coverage depends on CI with `dockerTest=true`.
<!--
1. Title: [#<issue>] <type>(<scope>): <subject>
Examples:
- "[#123] feat(operator): Support xxx"
- "[#233] fix: Check null before access result in xxx"
- "[MINOR] refactor: Fix typo in variable name"
- "[MINOR] docs: Fix typo in README"
- "[#255] test: Fix flaky test NameOfTheTest"
Reference: https://www.conventionalcommits.org/en/v1.0.0/
2. If the PR is unfinished, please mark this PR as draft.
-->
---------
Co-authored-by: lijie <[email protected]>
---
.../relational/mapper/TableVersionMapper.java | 16 +
.../mapper/TableVersionSQLProviderFactory.java | 14 +
.../provider/base/TableVersionBaseSQLProvider.java | 57 ++
.../postgresql/TableVersionPostgreSQLProvider.java | 42 +
.../relational/service/CatalogMetaService.java | 5 +
.../relational/service/FunctionMetaService.java | 76 +-
.../relational/service/MetalakeMetaService.java | 5 +
.../relational/service/SchemaMetaService.java | 30 +-
.../relational/service/TableMetaService.java | 106 ++-
.../relational/service/ViewMetaService.java | 72 +-
.../relational/service/TestCatalogMetaService.java | 68 ++
.../service/TestMetalakeMetaService.java | 73 ++
.../relational/service/TestSchemaMetaService.java | 906 +++++++++++++++++++++
13 files changed, 1393 insertions(+), 77 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionMapper.java
index 16f1ad7a34..af27c59d7d 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionMapper.java
@@ -19,6 +19,7 @@
package org.apache.gravitino.storage.relational.mapper;
+import java.util.List;
import org.apache.gravitino.storage.relational.po.TablePO;
import org.apache.ibatis.annotations.DeleteProvider;
import org.apache.ibatis.annotations.InsertProvider;
@@ -42,6 +43,21 @@ public interface TableVersionMapper {
void softDeleteTableVersionByTableIdAndVersion(
@Param("tableId") Long tableId, @Param("version") Long version);
+ @UpdateProvider(
+ type = TableVersionSQLProviderFactory.class,
+ method = "softDeleteTableVersionsBySchemaIds")
+ Integer softDeleteTableVersionsBySchemaIds(@Param("schemaIds") List<Long>
schemaIds);
+
+ @UpdateProvider(
+ type = TableVersionSQLProviderFactory.class,
+ method = "softDeleteTableVersionsByCatalogId")
+ Integer softDeleteTableVersionsByCatalogId(@Param("catalogId") Long
catalogId);
+
+ @UpdateProvider(
+ type = TableVersionSQLProviderFactory.class,
+ method = "softDeleteTableVersionsByMetalakeId")
+ Integer softDeleteTableVersionsByMetalakeId(@Param("metalakeId") Long
metalakeId);
+
@DeleteProvider(
type = TableVersionSQLProviderFactory.class,
method = "deleteTableVersionByLegacyTimeline")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionSQLProviderFactory.java
index 4c518ef4bd..b0d2315a0a 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableVersionSQLProviderFactory.java
@@ -20,6 +20,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.TableVersionBaseSQLProvider;
@@ -65,6 +66,19 @@ public class TableVersionSQLProviderFactory {
return getProvider().softDeleteTableVersionByTableIdAndVersion(tableId,
version);
}
+ public static String softDeleteTableVersionsBySchemaIds(
+ @Param("schemaIds") List<Long> schemaIds) {
+ return getProvider().softDeleteTableVersionsBySchemaIds(schemaIds);
+ }
+
+ public static String softDeleteTableVersionsByCatalogId(@Param("catalogId")
Long catalogId) {
+ return getProvider().softDeleteTableVersionsByCatalogId(catalogId);
+ }
+
+ public static String
softDeleteTableVersionsByMetalakeId(@Param("metalakeId") Long metalakeId) {
+ return getProvider().softDeleteTableVersionsByMetalakeId(metalakeId);
+ }
+
public static String deleteTableVersionByLegacyTimeline(
@Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
return getProvider().deleteTableVersionByLegacyTimeline(legacyTimeline,
limit);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableVersionBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableVersionBaseSQLProvider.java
index 2108a16a08..11dcc3a1be 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableVersionBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableVersionBaseSQLProvider.java
@@ -21,6 +21,8 @@ package
org.apache.gravitino.storage.relational.mapper.provider.base;
import static
org.apache.gravitino.storage.relational.mapper.TableVersionMapper.TABLE_NAME;
+import java.util.List;
+import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
import org.apache.gravitino.storage.relational.mapper.provider.DatabaseTimeSQL;
import org.apache.gravitino.storage.relational.po.TablePO;
import org.apache.ibatis.annotations.Param;
@@ -86,6 +88,61 @@ public class TableVersionBaseSQLProvider {
+ " WHERE table_id = #{tableId} AND version = #{version} AND
deleted_at = 0";
}
+ /**
+ * Soft-deletes all active table version rows whose parent table belongs to
one of the given
+ * schema IDs. The table_version_info table has no schema_id column, so a
sub-query joins
+ * table_meta to find the matching table_ids. The sub-query intentionally
does not filter on
+ * table_meta.deleted_at because this method runs after
softDeleteTableMetasBySchemaIds within the
+ * same transaction; at that point the table_meta rows already have
deleted_at set.
+ */
+ public String softDeleteTableVersionsBySchemaIds(@Param("schemaIds")
List<Long> schemaIds) {
+ return "<script>"
+ + "UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = "
+ + DatabaseTimeSQL.MYSQL
+ + " WHERE table_id IN (SELECT table_id FROM "
+ + TableMetaMapper.TABLE_NAME
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}"
+ + "</foreach>"
+ + ")) AND deleted_at = 0"
+ + "</script>";
+ }
+
+ /**
+ * Soft-deletes active table version rows whose parent table belongs to the
given catalog.
+ * table_version_info has no catalog_id column, so a sub-query on table_meta
finds the table_ids.
+ * Runs after softDeleteTableMetasByCatalogId in the same transaction, so
the sub-query does not
+ * filter table_meta.deleted_at.
+ */
+ public String softDeleteTableVersionsByCatalogId(@Param("catalogId") Long
catalogId) {
+ return "UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = "
+ + DatabaseTimeSQL.MYSQL
+ + " WHERE table_id IN (SELECT table_id FROM "
+ + TableMetaMapper.TABLE_NAME
+ + " WHERE catalog_id = #{catalogId}) AND deleted_at = 0";
+ }
+
+ /**
+ * Soft-deletes active table version rows whose parent table belongs to the
given metalake.
+ * table_version_info has no metalake_id column, so a sub-query on
table_meta finds the table_ids.
+ * Runs after softDeleteTableMetasByMetalakeId in the same transaction, so
the sub-query does not
+ * filter table_meta.deleted_at.
+ */
+ public String softDeleteTableVersionsByMetalakeId(@Param("metalakeId") Long
metalakeId) {
+ return "UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = "
+ + DatabaseTimeSQL.MYSQL
+ + " WHERE table_id IN (SELECT table_id FROM "
+ + TableMetaMapper.TABLE_NAME
+ + " WHERE metalake_id = #{metalakeId}) AND deleted_at = 0";
+ }
+
public String deleteTableVersionByLegacyTimeline(
@Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
return "DELETE FROM "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableVersionPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableVersionPostgreSQLProvider.java
index bc462c3839..1c77b843f8 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableVersionPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableVersionPostgreSQLProvider.java
@@ -21,6 +21,9 @@ package
org.apache.gravitino.storage.relational.mapper.provider.postgresql;
import static
org.apache.gravitino.storage.relational.mapper.TableVersionMapper.TABLE_NAME;
+import java.util.List;
+import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.provider.DatabaseTimeSQL;
import
org.apache.gravitino.storage.relational.mapper.provider.base.TableVersionBaseSQLProvider;
import org.apache.gravitino.storage.relational.po.TablePO;
import org.apache.ibatis.annotations.Param;
@@ -66,6 +69,45 @@ public class TableVersionPostgreSQLProvider extends
TableVersionBaseSQLProvider
+ " WHERE table_id = #{tableId} AND version = #{version} AND
deleted_at = 0";
}
+ @Override
+ public String softDeleteTableVersionsBySchemaIds(@Param("schemaIds")
List<Long> schemaIds) {
+ return "<script>"
+ + "UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = "
+ + DatabaseTimeSQL.POSTGRESQL
+ + " WHERE table_id IN (SELECT table_id FROM "
+ + TableMetaMapper.TABLE_NAME
+ + " WHERE schema_id IN ("
+ + "<foreach collection='schemaIds' item='schemaId' separator=','>"
+ + "#{schemaId}"
+ + "</foreach>"
+ + ")) AND deleted_at = 0"
+ + "</script>";
+ }
+
+ @Override
+ public String softDeleteTableVersionsByCatalogId(@Param("catalogId") Long
catalogId) {
+ return "UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = "
+ + DatabaseTimeSQL.POSTGRESQL
+ + " WHERE table_id IN (SELECT table_id FROM "
+ + TableMetaMapper.TABLE_NAME
+ + " WHERE catalog_id = #{catalogId}) AND deleted_at = 0";
+ }
+
+ @Override
+ public String softDeleteTableVersionsByMetalakeId(@Param("metalakeId") Long
metalakeId) {
+ return "UPDATE "
+ + TABLE_NAME
+ + " SET deleted_at = "
+ + DatabaseTimeSQL.POSTGRESQL
+ + " WHERE table_id IN (SELECT table_id FROM "
+ + TableMetaMapper.TABLE_NAME
+ + " WHERE metalake_id = #{metalakeId}) AND deleted_at = 0";
+ }
+
@Override
public String deleteTableVersionByLegacyTimeline(
@Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
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 fdf966321f..c80fe23411 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
@@ -51,6 +51,7 @@ import
org.apache.gravitino.storage.relational.mapper.SecurableObjectMapper;
import org.apache.gravitino.storage.relational.mapper.StatisticMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TableColumnMapper;
import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.TableVersionMapper;
import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
@@ -290,6 +291,10 @@ public class CatalogMetaService {
SessionUtils.doWithoutCommit(
TableMetaMapper.class,
mapper -> mapper.softDeleteTableMetasByCatalogId(catalogId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableVersionMapper.class,
+ mapper ->
mapper.softDeleteTableVersionsByCatalogId(catalogId)),
() ->
SessionUtils.doWithoutCommit(
TableColumnMapper.class,
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
index 5719534ae3..cf53adc5c6 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
@@ -267,28 +267,60 @@ public class FunctionMetaService {
try {
FunctionPO newFunctionPO =
updateFunctionPO(oldFunctionPO, newEntity, newSchemaId,
newCatalogId, newMetalakeId);
- SchemaMetaService.getInstance()
- .doWithSchemaWriteLock(
- newEntity.nameIdentifier(),
- newSchemaId,
- newCatalogId,
- newMetalakeId,
- () -> {
- // function_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(
- FunctionMetaMapper.class,
- mapper -> ops.updatePO(mapper, newFunctionPO,
oldFunctionPO));
- if (updated == 0) {
- throw functionWriteFailure(identifier, oldFunctionPO);
- }
- },
- () ->
- SessionUtils.doWithoutCommit(
- FunctionVersionMetaMapper.class,
- mapper ->
-
mapper.insertFunctionVersionMeta(newFunctionPO.functionVersionPO())));
+ if (isSchemaChanged) {
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ SchemaMetaService.getInstance()
+ .lockCatalogForEntityWrite(
+ oldFunctionEntity.nameIdentifier(),
+ oldFunctionPO.catalogId(),
+ oldFunctionPO.metalakeId()),
+ () ->
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ oldFunctionEntity.nameIdentifier(),
+ oldFunctionPO.schemaId(),
+ oldFunctionPO.catalogId(),
+ oldFunctionPO.metalakeId()),
+ () ->
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ newEntity.nameIdentifier(), newSchemaId, newCatalogId,
newMetalakeId),
+ () -> {
+ int updated =
+ SessionUtils.getWithoutCommit(
+ FunctionMetaMapper.class,
+ mapper -> ops.updatePO(mapper, newFunctionPO,
oldFunctionPO));
+ if (updated == 0) {
+ throw functionWriteFailure(identifier, oldFunctionPO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ FunctionVersionMetaMapper.class,
+ mapper ->
mapper.insertFunctionVersionMeta(newFunctionPO.functionVersionPO())));
+ } else {
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ newEntity.nameIdentifier(),
+ newSchemaId,
+ newCatalogId,
+ newMetalakeId,
+ () -> {
+ int updated =
+ SessionUtils.getWithoutCommit(
+ FunctionMetaMapper.class,
+ mapper -> ops.updatePO(mapper, newFunctionPO,
oldFunctionPO));
+ if (updated == 0) {
+ throw functionWriteFailure(identifier, oldFunctionPO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ FunctionVersionMetaMapper.class,
+ mapper ->
+
mapper.insertFunctionVersionMeta(newFunctionPO.functionVersionPO())));
+ }
return newEntity;
} catch (RuntimeException re) {
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 4a2a0faa29..ee811c7d44 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
@@ -57,6 +57,7 @@ import
org.apache.gravitino.storage.relational.mapper.SecurableObjectMapper;
import org.apache.gravitino.storage.relational.mapper.StatisticMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TableColumnMapper;
import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.TableVersionMapper;
import org.apache.gravitino.storage.relational.mapper.TagMetaMapper;
import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
@@ -229,6 +230,10 @@ public class MetalakeMetaService {
SessionUtils.doWithoutCommit(
TableMetaMapper.class,
mapper ->
mapper.softDeleteTableMetasByMetalakeId(metalakeId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableVersionMapper.class,
+ mapper ->
mapper.softDeleteTableVersionsByMetalakeId(metalakeId)),
() ->
SessionUtils.doWithoutCommit(
TableColumnMapper.class,
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 5e03b3d860..671af53f1c 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
@@ -58,6 +58,7 @@ import
org.apache.gravitino.storage.relational.mapper.SecurableObjectMapper;
import org.apache.gravitino.storage.relational.mapper.StatisticMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TableColumnMapper;
import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.TableVersionMapper;
import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
@@ -295,6 +296,10 @@ public class SchemaMetaService {
SessionUtils.doWithoutCommit(
TableMetaMapper.class,
mapper ->
mapper.softDeleteTableMetasBySchemaIds(schemaIds.get())),
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableVersionMapper.class,
+ mapper ->
mapper.softDeleteTableVersionsBySchemaIds(schemaIds.get())),
() ->
SessionUtils.doWithoutCommit(
TableColumnMapper.class,
@@ -519,7 +524,30 @@ public class SchemaMetaService {
SessionUtils.doMultipleWithCommit(transactionOperations);
}
- private void lockSchemaForEntityWrite(
+ /**
+ * Takes a shared lock on the parent catalog row before a cross-schema child
move.
+ *
+ * <p>A cascade schema delete holds an exclusive catalog lock before any
schema lock. Taking the
+ * same shared catalog lock here first ensures that a cross-schema child
move blocks the cascade
+ * delete until both schema locks are acquired, and a cascade delete blocks
new cross-schema moves
+ * until it finishes. Without this catalog fence, a move that holds schema A
could deadlock
+ * against a cascade that holds the catalog and is waiting for schema A.
+ */
+ void lockCatalogForEntityWrite(NameIdentifier entityIdentifier, Long
catalogId, Long metalakeId) {
+ String catalogName = entityIdentifier.namespace().level(1);
+ OccWriteSupport.lockParentForChildWrite(
+ catalogName,
+ Entity.EntityType.CATALOG,
+ () ->
+ SessionUtils.getWithoutCommit(
+ CatalogMetaMapper.class, mapper ->
mapper.selectCatalogMetaByIdForShare(catalogId)),
+ null,
+ current ->
+ Objects.equals(current.getCatalogName(), catalogName)
+ && Objects.equals(current.getMetalakeId(), metalakeId));
+ }
+
+ void lockSchemaForEntityWrite(
NameIdentifier entityIdentifier,
Long observedSchemaId,
Long observedCatalogId,
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
index e0da487bbd..ecbd87aa4a 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
@@ -219,41 +219,79 @@ public class TableMetaService {
POConverters.updateTablePOWithVersionAndSchemaId(oldTablePO,
newTableEntity, newSchemaId);
try {
- // For a cross-schema rename, the new schema is the parent that must
remain alive. For a
- // regular update, newSchemaId is the existing parent, so the same entry
point covers both.
- SchemaMetaService.getInstance()
- .doWithSchemaWriteLock(
- newTableEntity.nameIdentifier(),
- newSchemaId,
- oldTablePO.getCatalogId(),
- oldTablePO.getMetalakeId(),
- () -> {
- // current_version is the table's OCC token. A zero-row update
means another
- // writer won, so stop before touching version history or
columns.
- int updated =
- SessionUtils.getWithoutCommit(
- TableMetaMapper.class,
- mapper -> ops.updatePO(mapper, newTablePO,
oldTablePO));
- if (updated == 0) {
- throw tableWriteFailure(identifier, oldTablePO);
- }
- },
- () ->
- SessionUtils.doWithoutCommit(
- TableVersionMapper.class,
- mapper -> {
- // Only the CAS winner can reach this step, so it is
safe to replace the
- // details stored under the next table version.
- mapper.softDeleteTableVersionByTableIdAndVersion(
- oldTablePO.getTableId(),
oldTablePO.getCurrentVersion());
-
mapper.insertTableVersionOnDuplicateKeyUpdate(newTablePO);
- }),
- () -> {
- // A column failure rolls back the table row and version row
in the same
- // transaction.
+ if (isSchemaChanged) {
+ // A cross-schema move touches two schemas. Lock the catalog first so
a cascade
+ // schema-delete (which also locks the catalog exclusively) cannot
slip between the
+ // two schema locks and cause a deadlock.
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ SchemaMetaService.getInstance()
+ .lockCatalogForEntityWrite(
+ oldTableEntity.nameIdentifier(),
+ oldTablePO.getCatalogId(),
+ oldTablePO.getMetalakeId()),
+ () ->
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ oldTableEntity.nameIdentifier(),
+ oldTablePO.getSchemaId(),
+ oldTablePO.getCatalogId(),
+ oldTablePO.getMetalakeId()),
+ () ->
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ newTableEntity.nameIdentifier(),
+ newSchemaId,
+ oldTablePO.getCatalogId(),
+ oldTablePO.getMetalakeId()),
+ () -> {
+ int updated =
+ SessionUtils.getWithoutCommit(
+ TableMetaMapper.class,
+ mapper -> ops.updatePO(mapper, newTablePO, oldTablePO));
+ if (updated == 0) {
+ throw tableWriteFailure(identifier, oldTablePO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableVersionMapper.class,
+ mapper -> {
+ mapper.softDeleteTableVersionByTableIdAndVersion(
+ oldTablePO.getTableId(),
oldTablePO.getCurrentVersion());
+
mapper.insertTableVersionOnDuplicateKeyUpdate(newTablePO);
+ }),
+ () ->
TableColumnMetaService.getInstance()
- .updateColumnPOsFromTableDiff(oldTableEntity,
newTableEntity, newTablePO);
- });
+ .updateColumnPOsFromTableDiff(oldTableEntity,
newTableEntity, newTablePO));
+ } else {
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ newTableEntity.nameIdentifier(),
+ newSchemaId,
+ oldTablePO.getCatalogId(),
+ oldTablePO.getMetalakeId(),
+ () -> {
+ int updated =
+ SessionUtils.getWithoutCommit(
+ TableMetaMapper.class,
+ mapper -> ops.updatePO(mapper, newTablePO,
oldTablePO));
+ if (updated == 0) {
+ throw tableWriteFailure(identifier, oldTablePO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableVersionMapper.class,
+ mapper -> {
+ mapper.softDeleteTableVersionByTableIdAndVersion(
+ oldTablePO.getTableId(),
oldTablePO.getCurrentVersion());
+
mapper.insertTableVersionOnDuplicateKeyUpdate(newTablePO);
+ }),
+ () ->
+ TableColumnMetaService.getInstance()
+ .updateColumnPOsFromTableDiff(oldTableEntity,
newTableEntity, newTablePO));
+ }
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
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 d5b2f8026f..afdc379db8 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
@@ -159,26 +159,58 @@ public class ViewMetaService {
try {
ViewPO newViewPO = updateViewPO(oldViewPO, newEntity);
- SchemaMetaService.getInstance()
- .doWithSchemaWriteLock(
- 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 (updated == 0) {
- throw viewWriteFailure(ident, oldViewPO);
- }
- },
- () ->
- SessionUtils.doWithoutCommit(
- ViewVersionInfoMapper.class,
- mapper ->
mapper.insertViewVersionInfo(newViewPO.getViewVersionInfoPO())));
+ if (isSchemaChanged) {
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ SchemaMetaService.getInstance()
+ .lockCatalogForEntityWrite(
+ oldViewEntity.nameIdentifier(),
+ oldViewPO.getCatalogId(),
+ oldViewPO.getMetalakeId()),
+ () ->
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ oldViewEntity.nameIdentifier(),
+ oldViewPO.getSchemaId(),
+ oldViewPO.getCatalogId(),
+ oldViewPO.getMetalakeId()),
+ () ->
+ SchemaMetaService.getInstance()
+ .lockSchemaForEntityWrite(
+ newEntity.nameIdentifier(), newSchemaId, newCatalogId,
newMetalakeId),
+ () -> {
+ int updated =
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class, mapper -> ops.updatePO(mapper,
newViewPO, oldViewPO));
+ if (updated == 0) {
+ throw viewWriteFailure(ident, oldViewPO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.insertViewVersionInfo(newViewPO.getViewVersionInfoPO())));
+ } else {
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ newEntity.nameIdentifier(),
+ newSchemaId,
+ newCatalogId,
+ newMetalakeId,
+ () -> {
+ int updated =
+ SessionUtils.getWithoutCommit(
+ 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) {
ExceptionUtils.checkSQLException(
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestCatalogMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestCatalogMetaService.java
index cbc598903e..2a6bc7c441 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestCatalogMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestCatalogMetaService.java
@@ -573,6 +573,55 @@ public class TestCatalogMetaService extends
TestJDBCBackend {
assertEquals(0, countActiveTagRelForMetadataObject(function.id(),
"FUNCTION"));
}
+ @TestTemplate
+ public void testDeleteCatalogCascadeRemovesTableVersions() throws
IOException {
+ CatalogEntity catalog =
+ createCatalog(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofCatalog(metalakeName),
+ "catalog_with_table_versions",
+ auditInfo);
+ backend.insert(catalog, false);
+
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(metalakeName, catalog.name()),
+ "schema_with_table_versions",
+ AUDIT_INFO);
+ backend.insert(schema, false);
+
+ Namespace objectNamespace = Namespace.of(metalakeName, catalog.name(),
schema.name());
+ ColumnEntity column =
+ ColumnEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("column_with_version")
+ .withPosition(0)
+ .withAutoIncrement(false)
+ .withNullable(false)
+ .withDataType(Types.IntegerType.get())
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TableEntity table =
+ TableEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("table_with_version")
+ .withNamespace(objectNamespace)
+ .withColumns(List.of(column))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TableMetaService.getInstance().insertTable(table, false);
+
+ // The insert wrote an active version row for the table.
+ assertTrue(countActiveTableVersionRows(table.id()) > 0);
+
+
assertTrue(CatalogMetaService.getInstance().deleteCatalog(catalog.nameIdentifier(),
true));
+
+ // The catalog cascade must soft-delete the version rows too; otherwise
they keep
+ // deleted_at = 0 forever and never become eligible for legacy-timeline
cleanup.
+ assertEquals(0, countActiveTableVersionRows(table.id()));
+ }
+
private List<Throwable> insertCatalogsConcurrently(CatalogEntity first,
CatalogEntity second)
throws Exception {
ExecutorService executor = Executors.newFixedThreadPool(2);
@@ -656,4 +705,23 @@ public class TestCatalogMetaService extends
TestJDBCBackend {
throw new RuntimeException("SQL execution failed", e);
}
}
+
+ private int countActiveTableVersionRows(Long tableId) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet rs =
+ statement.executeQuery(
+ String.format(
+ "SELECT count(*) FROM table_version_info WHERE table_id =
%d AND deleted_at = 0",
+ tableId))) {
+ if (rs.next()) {
+ return rs.getInt(1);
+ }
+ return 0;
+ } catch (SQLException e) {
+ throw new RuntimeException("SQL execution failed", e);
+ }
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestMetalakeMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestMetalakeMetaService.java
index 952416c1cf..ed560ee704 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestMetalakeMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestMetalakeMetaService.java
@@ -18,11 +18,16 @@
*/
package org.apache.gravitino.storage.relational.service;
+import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.io.IOException;
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
import java.time.Instant;
import java.util.List;
import java.util.concurrent.CountDownLatch;
@@ -33,22 +38,28 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import org.apache.gravitino.Entity;
import org.apache.gravitino.EntityAlreadyExistsException;
+import org.apache.gravitino.Namespace;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.exceptions.NonEmptyEntityException;
import org.apache.gravitino.exceptions.OptimisticLockException;
import org.apache.gravitino.meta.BaseMetalake;
import org.apache.gravitino.meta.CatalogEntity;
+import org.apache.gravitino.meta.ColumnEntity;
import org.apache.gravitino.meta.SchemaEntity;
import org.apache.gravitino.meta.SchemaVersion;
+import org.apache.gravitino.meta.TableEntity;
+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.MetalakeMetaMapper;
import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
import org.apache.gravitino.storage.relational.po.MetalakePO;
import org.apache.gravitino.storage.relational.po.SchemaPO;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
import org.apache.gravitino.storage.relational.utils.POConverters;
import org.apache.gravitino.storage.relational.utils.SessionUtils;
import org.apache.gravitino.utils.NamespaceUtil;
+import org.apache.ibatis.session.SqlSession;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.TestTemplate;
import org.mockito.Mockito;
@@ -415,6 +426,49 @@ public class TestMetalakeMetaService extends
TestJDBCBackend {
assertTrue(backend.exists(metalake.nameIdentifier(),
Entity.EntityType.METALAKE));
}
+ @TestTemplate
+ public void testDeleteMetalakeCascadeRemovesTableVersions() throws
IOException {
+ BaseMetalake metalake = createAndInsertMakeLake(METALAKE_NAME);
+ CatalogEntity catalog = createAndInsertCatalog(METALAKE_NAME,
"catalog_with_table_versions");
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(METALAKE_NAME, catalog.name()),
+ "schema_with_table_versions",
+ AUDIT_INFO);
+ backend.insert(schema, false);
+
+ Namespace objectNamespace = Namespace.of(METALAKE_NAME, catalog.name(),
schema.name());
+ ColumnEntity column =
+ ColumnEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("column_with_version")
+ .withPosition(0)
+ .withAutoIncrement(false)
+ .withNullable(false)
+ .withDataType(Types.IntegerType.get())
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TableEntity table =
+ TableEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("table_with_version")
+ .withNamespace(objectNamespace)
+ .withColumns(List.of(column))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TableMetaService.getInstance().insertTable(table, false);
+
+ // The insert wrote an active version row for the table.
+ assertTrue(countActiveTableVersionRows(table.id()) > 0);
+
+
assertTrue(MetalakeMetaService.getInstance().deleteMetalake(metalake.nameIdentifier(),
true));
+
+ // The metalake cascade must soft-delete the version rows too; otherwise
they keep
+ // deleted_at = 0 forever and never become eligible for legacy-timeline
cleanup.
+ assertEquals(0, countActiveTableVersionRows(table.id()));
+ }
+
@TestTemplate
public void testMetaLifeCycleFromCreationToDeletion() throws IOException {
// meta data creation
@@ -442,4 +496,23 @@ public class TestMetalakeMetaService extends
TestJDBCBackend {
}
assertFalse(legacyRecordExistsInDB(metalake.id(),
Entity.EntityType.METALAKE));
}
+
+ private int countActiveTableVersionRows(Long tableId) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet rs =
+ statement.executeQuery(
+ String.format(
+ "SELECT count(*) FROM table_version_info WHERE table_id =
%d AND deleted_at = 0",
+ tableId))) {
+ if (rs.next()) {
+ return rs.getInt(1);
+ }
+ return 0;
+ } catch (SQLException e) {
+ throw new RuntimeException("SQL execution failed", e);
+ }
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
index 4ea5eea35c..60a0dded70 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
@@ -22,16 +22,20 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import com.google.common.collect.ImmutableMap;
import java.io.IOException;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.time.Instant;
+import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.EnumMap;
import java.util.List;
import java.util.Locale;
+import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
@@ -53,13 +57,16 @@ import org.apache.gravitino.meta.ColumnEntity;
import org.apache.gravitino.meta.FilesetEntity;
import org.apache.gravitino.meta.FunctionEntity;
import org.apache.gravitino.meta.ModelEntity;
+import org.apache.gravitino.meta.ModelVersionEntity;
import org.apache.gravitino.meta.SchemaEntity;
import org.apache.gravitino.meta.TableEntity;
import org.apache.gravitino.meta.TagEntity;
import org.apache.gravitino.meta.TopicEntity;
import org.apache.gravitino.meta.ViewEntity;
+import org.apache.gravitino.model.ModelVersion;
import org.apache.gravitino.rel.types.Types;
import org.apache.gravitino.storage.RandomIdGenerator;
+import org.apache.gravitino.storage.relational.RelationalBackend;
import org.apache.gravitino.storage.relational.TestJDBCBackend;
import org.apache.gravitino.storage.relational.mapper.CatalogMetaMapper;
import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
@@ -777,6 +784,530 @@ public class TestSchemaMetaService extends
TestJDBCBackend {
NameIdentifier.of(metalakeName, catalogName, "anc_a"),
Entity.EntityType.SCHEMA));
}
+ @TestTemplate
+ public void testSchemaChildUpdateServicesWaitForConcurrentSchemaDelete()
throws Exception {
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+
+ List<SchemaChildUpdateCase> childCases = schemaChildUpdateCases();
+
+ for (int index = 0; index < childCases.size(); index++) {
+ SchemaChildUpdateCase childCase = childCases.get(index);
+ String schemaName =
+ "schema_for_update_lock_" +
childCase.entityType.name().toLowerCase(Locale.ROOT);
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ schemaName,
+ AUDIT_INFO);
+ backend.insert(schema, false);
+
+ Namespace schemaNamespace = Namespace.of(metalakeName, catalogName,
schemaName);
+ String childName = "child_" +
childCase.entityType.name().toLowerCase(Locale.ROOT);
+ Object[] ref = childCase.createChild.run(schemaNamespace, childName,
backend);
+ NameIdentifier childIdent = (NameIdentifier) ref[0];
+ Long childId = (Long) ref[1];
+
+ assertChildUpdateBlocksOnConcurrentSchemaDelete(schema, childIdent,
childCase, childId);
+ }
+ }
+
+ private void assertChildUpdateBlocksOnConcurrentSchemaDelete(
+ SchemaEntity schema, NameIdentifier childIdent, SchemaChildUpdateCase
childCase, Long childId)
+ throws Exception {
+ SchemaPO observedSchemaPO =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class, mapper ->
mapper.selectSchemaMetaById(schema.id()));
+
+ CountDownLatch schemaDeleteLocked = new CountDownLatch(1);
+ CountDownLatch allowDeleteCommit = new CountDownLatch(1);
+ CountDownLatch updateStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+
+ // Thread 1: soft-delete the schema row and hold the transaction open. The
+ // uncommitted UPDATE holds an exclusive row lock on the schema_meta row.
+ Future<Throwable> deleteResult =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ int deleted =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+
mapper.softDeleteSchemaMetaBySchemaIdAndVersion(
+ observedSchemaPO.getSchemaId(),
+ observedSchemaPO.getCurrentVersion()));
+ Assertions.assertEquals(1, deleted);
+ schemaDeleteLocked.countDown();
+ try {
+ assertTrue(allowDeleteCommit.await(30,
TimeUnit.SECONDS));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ try {
+ assertTrue(schemaDeleteLocked.await(30, TimeUnit.SECONDS));
+
+ // Thread 2: attempt to update the child entity. The update's
doWithSchemaWriteLock
+ // must block on the schema row held by Thread 1's uncommitted
soft-delete. If the
+ // update completes within 500 ms, the schema lock was not acquired by
the update path.
+ Future<Throwable> updateResult =
+ executor.submit(
+ () -> {
+ try {
+ // An unlocked read first proves the worker is actually
running, so the
+ // timeout below measures lock contention, not scheduling
latency.
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class, mapper ->
mapper.selectSchemaMetaById(schema.id()));
+ updateStarted.countDown();
+ childCase.update.run(childIdent);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ assertTrue(updateStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> updateResult.get(500,
TimeUnit.MILLISECONDS));
+
+ // Commit Thread 1's schema-row soft-delete. Child-row cleanup is not
part of that
+ // transaction; it is covered by
testCascadeDeleteLeavesNoOrphanVersionRows.
+ allowDeleteCommit.countDown();
+ Assertions.assertNull(
+ deleteResult.get(30, TimeUnit.SECONDS), () -> "Schema soft-delete
transaction failed");
+
+ // The update should now fail because the schema row is committed as
deleted.
+ Throwable updateFailure = updateResult.get(30, TimeUnit.SECONDS);
+ Throwable root = updateFailure;
+ while (root != null && !(root instanceof NoSuchEntityException)) {
+ root = root.getCause();
+ }
+ Assertions.assertNotNull(
+ root,
+ () ->
+ "Expected NoSuchEntityException, but got "
+ + updateFailure.getClass().getSimpleName()
+ + ": "
+ + updateFailure.getMessage());
+
+ // Only the schema row was soft-deleted above; the child rows must still
be active.
+ assertActiveChildRowsRemain(childCase, childId);
+ assertFalse(backend.exists(schema.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ } finally {
+ allowDeleteCommit.countDown();
+ executor.shutdownNow();
+ }
+ }
+
+ @TestTemplate
+ public void testCascadeDeleteLeavesNoOrphanVersionRows() throws Exception {
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+
+ List<SchemaChildUpdateCase> childCases = schemaChildUpdateCases();
+
+ for (int index = 0; index < childCases.size(); index++) {
+ SchemaChildUpdateCase childCase = childCases.get(index);
+ String schemaName =
+ "schema_for_orphan_test_" +
childCase.entityType.name().toLowerCase(Locale.ROOT);
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ schemaName,
+ AUDIT_INFO);
+ SchemaMetaService.getInstance().insertSchema(schema, false);
+
+ Namespace childNamespace = Namespace.of(metalakeName, catalogName,
schemaName);
+ String childName = "child_" +
childCase.entityType.name().toLowerCase(Locale.ROOT);
+ Object[] ref = childCase.createChild.run(childNamespace, childName,
backend);
+ NameIdentifier childIdent = (NameIdentifier) ref[0];
+ Long childId = (Long) ref[1];
+
+ SchemaMetaService.getInstance().deleteSchema(schema.nameIdentifier(),
true);
+
+ Throwable updateFailure =
+ assertThrows(
+ Exception.class,
+ () -> childCase.update.run(childIdent),
+ () -> "Update should have failed for " + childCase.entityType);
+
+ Throwable root = updateFailure;
+ while (root != null && !(root instanceof NoSuchEntityException)) {
+ root = root.getCause();
+ }
+ Assertions.assertNotNull(
+ root,
+ () ->
+ "Expected NoSuchEntityException for "
+ + childCase.entityType
+ + ", but got "
+ + updateFailure.getClass().getSimpleName()
+ + ": "
+ + updateFailure.getMessage());
+
+ assertFalse(backend.exists(schema.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ assertFalse(backend.exists(childIdent, childCase.entityType));
+ assertNoActiveChildRowsRemain(childCase, childId);
+ }
+ }
+
+ @TestTemplate
+ public void testSchemaChildUpdateLockBlocksConcurrentCascadeDelete() throws
Exception {
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+
+ List<SchemaChildUpdateCase> childCases = schemaChildUpdateCases();
+
+ for (int index = 0; index < childCases.size(); index++) {
+ SchemaChildUpdateCase childCase = childCases.get(index);
+ String schemaName =
+ "schema_update_first_" +
childCase.entityType.name().toLowerCase(Locale.ROOT);
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ schemaName,
+ AUDIT_INFO);
+ backend.insert(schema, false);
+
+ Namespace childNamespace = Namespace.of(metalakeName, catalogName,
schemaName);
+ String childName = "child_" +
childCase.entityType.name().toLowerCase(Locale.ROOT);
+ Object[] ref = childCase.createChild.run(childNamespace, childName,
backend);
+ NameIdentifier childIdent = (NameIdentifier) ref[0];
+ Long childId = (Long) ref[1];
+
+ assertSchemaUpdateLockBlocksCascadeDelete(schema, childIdent, childCase,
childId);
+ }
+ }
+
+ private void assertSchemaUpdateLockBlocksCascadeDelete(
+ SchemaEntity schema, NameIdentifier childIdent, SchemaChildUpdateCase
childCase, Long childId)
+ throws Exception {
+ CountDownLatch updateHoldsSchemaLock = new CountDownLatch(1);
+ CountDownLatch releaseUpdateLock = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+
+ // Thread 1: hold the schema row's shared lock, as a child update does in
+ // lockSchemaForEntityWrite, and keep the transaction open until released.
+ Future<Throwable> updateLockHolder =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ SchemaPO locked =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
mapper.selectSchemaMetaByIdForShare(schema.id()));
+ Assertions.assertNotNull(locked);
+ updateHoldsSchemaLock.countDown();
+ try {
+ assertTrue(releaseUpdateLock.await(30,
TimeUnit.SECONDS));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ try {
+ assertTrue(updateHoldsSchemaLock.await(30, TimeUnit.SECONDS));
+
+ // Thread 2: run the real cascade delete. Updating the schema row
conflicts with
+ // the shared lock above, so the delete must wait for Thread 1.
+ Future<Throwable> deleteResult =
+ executor.submit(
+ () -> {
+ try {
+
SchemaMetaService.getInstance().deleteSchema(schema.nameIdentifier(), true);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ // If the cascade finishes within 500 ms it never waited for the schema
lock.
+ assertThrows(TimeoutException.class, () -> deleteResult.get(500,
TimeUnit.MILLISECONDS));
+
+ releaseUpdateLock.countDown();
+ Assertions.assertNull(updateLockHolder.get(30, TimeUnit.SECONDS));
+ Assertions.assertNull(deleteResult.get(30, TimeUnit.SECONDS), () ->
"Cascade delete failed");
+
+ // The cascade completed after the lock release: no active child rows
remain.
+ assertFalse(backend.exists(schema.nameIdentifier(),
Entity.EntityType.SCHEMA));
+ assertFalse(backend.exists(childIdent, childCase.entityType));
+ assertNoActiveChildRowsRemain(childCase, childId);
+ } finally {
+ releaseUpdateLock.countDown();
+ executor.shutdownNow();
+ }
+ }
+
+ @TestTemplate
+ public void testCrossSchemaMoveAndHierarchicalCascadeDeleteLockOrder()
throws Exception {
+ createAndInsertMakeLake(metalakeName);
+ CatalogEntity catalog = createAndInsertCatalog(metalakeName, catalogName);
+ Long catalogId = catalog.id();
+
+ Map<Entity.EntityType, SchemaChildMove> moves = schemaChildMoves();
+
+ int index = 0;
+ for (SchemaChildUpdateCase moveCase : schemaChildUpdateCases()) {
+ SchemaChildMove move = moves.get(moveCase.entityType);
+ if (move == null) {
+ // Only table, view and function support cross-schema moves.
+ continue;
+ }
+ String type = moveCase.entityType.name().toLowerCase(Locale.ROOT);
+
+ // Fixed IDs with the child schema ID smaller than the parent's: the
pre-fix lock
+ // order (smaller schemaId first) deadlocked with the hierarchical
cascade delete
+ // under this assignment. Each direction uses fresh names/IDs because
direction 1
+ // soft-deletes its fixtures.
+ long parentSchemaIdA = 2000L + index;
+ long childSchemaIdA = 1000L + index;
+ String parentSchemaNameA = "move_parent_a_" + type;
+ String childSchemaNameA = parentSchemaNameA + ":move_child";
+
+ SchemaMetaService.getInstance()
+ .insertSchema(
+ createSchemaEntity(
+ parentSchemaIdA,
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ parentSchemaNameA,
+ AUDIT_INFO),
+ false);
+ SchemaMetaService.getInstance()
+ .insertSchema(
+ createSchemaEntity(
+ childSchemaIdA,
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ childSchemaNameA,
+ AUDIT_INFO),
+ false);
+
+ // The child leaf insert must not have recreated the ancestor with a new
ID.
+ Assertions.assertEquals(
+ parentSchemaIdA,
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
mapper.selectSchemaMetaByCatalogIdAndName(catalogId, parentSchemaNameA))
+ .getSchemaId());
+
+ Object[] ref =
+ moveCase.createChild.run(
+ Namespace.of(metalakeName, catalogName, childSchemaNameA),
+ "move_child_" + type,
+ backend);
+ NameIdentifier childIdent = (NameIdentifier) ref[0];
+ Long childId = (Long) ref[1];
+
+ // Direction 1: a held catalog shared lock blocks the cascade delete's
exclusive
+ // catalog lock.
+ assertCatalogSharedLockBlocksCascadeDelete(
+ catalogId, parentSchemaNameA, childIdent, moveCase, childId);
+
+ // Direction 2: a held catalog exclusive lock blocks the real
cross-schema move,
+ // whose first lock step is the catalog shared lock in
lockCatalogForEntityWrite.
+ long parentSchemaIdB = 4000L + index;
+ long childSchemaIdB = 3000L + index;
+ String parentSchemaNameB = "move_parent_b_" + type;
+ String childSchemaNameB = parentSchemaNameB + ":move_child";
+ SchemaMetaService.getInstance()
+ .insertSchema(
+ createSchemaEntity(
+ parentSchemaIdB,
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ parentSchemaNameB,
+ AUDIT_INFO),
+ false);
+ SchemaMetaService.getInstance()
+ .insertSchema(
+ createSchemaEntity(
+ childSchemaIdB,
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ childSchemaNameB,
+ AUDIT_INFO),
+ false);
+ Object[] ref2 =
+ moveCase.createChild.run(
+ Namespace.of(metalakeName, catalogName, childSchemaNameB),
+ "move_child_" + type,
+ backend);
+ NameIdentifier childIdent2 = (NameIdentifier) ref2[0];
+ assertCatalogExclusiveLockBlocksCrossSchemaMove(
+ catalogId, parentSchemaNameB, childIdent2, moveCase, move);
+ index++;
+ }
+ }
+
+ private void assertCatalogSharedLockBlocksCascadeDelete(
+ Long catalogId,
+ String parentSchemaName,
+ NameIdentifier childIdent,
+ SchemaChildUpdateCase moveCase,
+ Long childId)
+ throws Exception {
+ CountDownLatch moveHoldsCatalogLock = new CountDownLatch(1);
+ CountDownLatch releaseMoveLock = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+
+ // Thread 1: hold the catalog shared lock, as a cross-schema move does
first in
+ // lockCatalogForEntityWrite, until released.
+ Future<Throwable> moveLockHolder =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ CatalogPO locked =
+ SessionUtils.getWithoutCommit(
+ CatalogMetaMapper.class,
+ mapper ->
mapper.selectCatalogMetaByIdForShare(catalogId));
+ Assertions.assertNotNull(locked);
+ moveHoldsCatalogLock.countDown();
+ try {
+ assertTrue(releaseMoveLock.await(30,
TimeUnit.SECONDS));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ try {
+ assertTrue(moveHoldsCatalogLock.await(30, TimeUnit.SECONDS));
+
+ // Thread 2: run the real hierarchical cascade delete. Its catalog
exclusive lock
+ // conflicts with the shared lock above, so the delete must wait for
Thread 1.
+ NameIdentifier parentIdent = NameIdentifier.of(metalakeName,
catalogName, parentSchemaName);
+ Future<Throwable> deleteResult =
+ executor.submit(
+ () -> {
+ try {
+ SchemaMetaService.getInstance().deleteSchema(parentIdent,
true);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ // If the cascade finishes within 500 ms it never waited for the catalog
lock.
+ assertThrows(TimeoutException.class, () -> deleteResult.get(500,
TimeUnit.MILLISECONDS));
+
+ releaseMoveLock.countDown();
+ Assertions.assertNull(moveLockHolder.get(30, TimeUnit.SECONDS));
+ Assertions.assertNull(deleteResult.get(30, TimeUnit.SECONDS), () ->
"Cascade delete failed");
+
+ // The cascade removed the parent and child schemas and every child row.
+ assertFalse(backend.exists(parentIdent, Entity.EntityType.SCHEMA));
+ assertFalse(backend.exists(childIdent, moveCase.entityType));
+ assertNoActiveChildRowsRemain(moveCase, childId);
+ } finally {
+ releaseMoveLock.countDown();
+ executor.shutdownNow();
+ }
+ }
+
+ private void assertCatalogExclusiveLockBlocksCrossSchemaMove(
+ Long catalogId,
+ String parentSchemaName,
+ NameIdentifier childIdent,
+ SchemaChildUpdateCase moveCase,
+ SchemaChildMove move)
+ throws Exception {
+ CountDownLatch deleteHoldsCatalogLock = new CountDownLatch(1);
+ CountDownLatch releaseDeleteLock = new CountDownLatch(1);
+ CountDownLatch moveStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+
+ // Thread 1: hold the catalog exclusive lock, as a hierarchical cascade
delete does
+ // first in lockCatalogForSchemaDelete, until released.
+ Future<Throwable> deleteLockHolder =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ CatalogPO locked =
+ SessionUtils.getWithoutCommit(
+ CatalogMetaMapper.class,
+ mapper ->
mapper.selectCatalogMetaByIdForUpdate(catalogId));
+ Assertions.assertNotNull(locked);
+ deleteHoldsCatalogLock.countDown();
+ try {
+ assertTrue(releaseDeleteLock.await(30,
TimeUnit.SECONDS));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ try {
+ assertTrue(deleteHoldsCatalogLock.await(30, TimeUnit.SECONDS));
+
+ // Thread 2: run the real cross-schema move. Its first lock is the
catalog shared
+ // lock, which must wait for Thread 1. The pre-fix order (schema locks
first) would
+ // not touch the catalog row and would finish within 500 ms.
+ Future<Throwable> moveResult =
+ executor.submit(
+ () -> {
+ try {
+ // An unlocked read first proves the worker is actually
running, so the
+ // timeout below measures the lock wait, not scheduling
latency.
+ SessionUtils.getWithoutCommit(
+ CatalogMetaMapper.class, mapper ->
mapper.selectCatalogMetaById(catalogId));
+ moveStarted.countDown();
+ move.run(childIdent, parentSchemaName);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+
+ assertTrue(moveStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> moveResult.get(500,
TimeUnit.MILLISECONDS));
+
+ releaseDeleteLock.countDown();
+ Assertions.assertNull(deleteLockHolder.get(30, TimeUnit.SECONDS));
+ Assertions.assertNull(
+ moveResult.get(30, TimeUnit.SECONDS), () -> "Move failed after the
lock was released");
+
+ // The move completed: the entity now lives directly under the parent
schema.
+ NameIdentifier movedIdent =
+ NameIdentifier.of(
+ Namespace.of(metalakeName, catalogName, parentSchemaName),
childIdent.name());
+ assertTrue(backend.exists(movedIdent, moveCase.entityType));
+ } finally {
+ releaseDeleteLock.countDown();
+ executor.shutdownNow();
+ }
+ }
+
@TestTemplate
public void testOverlappingHierarchicalSchemaDeletesDoNotDeadlock() throws
Exception {
createAndInsertMakeLake(metalakeName);
@@ -1302,4 +1833,379 @@ public class TestSchemaMetaService extends
TestJDBCBackend {
this.write = write;
}
}
+
+ @FunctionalInterface
+ private interface SchemaChildUpdate {
+ void run(NameIdentifier childIdentifier) throws Exception;
+ }
+
+ private int countActiveRowsForEntity(Long entityId, String tableName, String
idColumnName) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet rs =
+ statement.executeQuery(
+ String.format(
+ "SELECT count(*) FROM %s WHERE %s = %d AND deleted_at = 0",
+ tableName, idColumnName, entityId))) {
+ if (rs.next()) {
+ return rs.getInt(1);
+ }
+ return 0;
+ } catch (SQLException e) {
+ throw new RuntimeException("SQL execution failed", e);
+ }
+ }
+
+ private void assertActiveChildRowsRemain(SchemaChildUpdateCase childCase,
Long childId) {
+ assertChildRowCounts(childCase, childId, true);
+ }
+
+ private void assertNoActiveChildRowsRemain(SchemaChildUpdateCase childCase,
Long childId) {
+ assertChildRowCounts(childCase, childId, false);
+ }
+
+ private void assertChildRowCounts(
+ SchemaChildUpdateCase childCase, Long childId, boolean expectPresent) {
+ List<String[]> tables = new ArrayList<>();
+ tables.add(new String[] {childCase.metaTable, childCase.idColumn});
+ if (childCase.versionTable != null) {
+ tables.add(new String[] {childCase.versionTable, childCase.idColumn});
+ }
+ tables.addAll(childCase.extraTables);
+ for (String[] table : tables) {
+ int rows = countActiveRowsForEntity(childId, table[0], table[1]);
+ if (expectPresent) {
+ assertTrue(
+ rows > 0,
+ () ->
+ "Expected active "
+ + table[0]
+ + " rows for "
+ + childCase.entityType
+ + " after schema-only soft-delete");
+ } else {
+ Assertions.assertEquals(
+ 0,
+ rows,
+ () ->
+ "Found active "
+ + table[0]
+ + " rows for "
+ + childCase.entityType
+ + " after cascade delete");
+ }
+ }
+ }
+
+ @FunctionalInterface
+ private interface SchemaChildCreator {
+ Object[] run(Namespace namespace, String name, RelationalBackend backend)
throws Exception;
+ }
+
+ private static class SchemaChildUpdateCase {
+ private final Entity.EntityType entityType;
+ private final String metaTable;
+ private final String idColumn;
+ private final String versionTable;
+ private final List<String[]> extraTables;
+ private final SchemaChildCreator createChild;
+ private final SchemaChildUpdate update;
+
+ private SchemaChildUpdateCase(
+ Entity.EntityType entityType,
+ String metaTable,
+ String idColumn,
+ String versionTable,
+ List<String[]> extraTables,
+ SchemaChildCreator createChild,
+ SchemaChildUpdate update) {
+ this.entityType = entityType;
+ this.metaTable = metaTable;
+ this.idColumn = idColumn;
+ this.versionTable = versionTable;
+ this.extraTables = extraTables;
+ this.createChild = createChild;
+ this.update = update;
+ }
+ }
+
+ private List<SchemaChildUpdateCase> schemaChildUpdateCases() {
+ return Arrays.asList(
+ new SchemaChildUpdateCase(
+ Entity.EntityType.TABLE,
+ "table_meta",
+ "table_id",
+ "table_version_info",
+ Collections.singletonList(new String[]
{"table_column_version_info", "table_id"}),
+ (namespace, name, bk) -> {
+ ColumnEntity column =
+ ColumnEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("col_" + name)
+ .withPosition(0)
+ .withAutoIncrement(false)
+ .withNullable(false)
+ .withDataType(Types.IntegerType.get())
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TableEntity e =
+ TableEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(name)
+ .withNamespace(namespace)
+ .withAuditInfo(AUDIT_INFO)
+ .withColumns(List.of(column))
+ .build();
+ bk.insert(e, false);
+ return new Object[] {e.nameIdentifier(), e.id()};
+ },
+ childIdent ->
+ TableMetaService.getInstance()
+ .updateTable(
+ childIdent,
+ entity -> {
+ TableEntity t = (TableEntity) entity;
+ return TableEntity.builder()
+ .withId(t.id())
+ .withName(t.name())
+ .withNamespace(t.namespace())
+ .withAuditInfo(t.auditInfo())
+ .withColumns(t.columns())
+ .withComment("updated comment")
+ .withProperties(t.properties())
+ .build();
+ })),
+ new SchemaChildUpdateCase(
+ Entity.EntityType.VIEW,
+ "view_meta",
+ "view_id",
+ "view_version_info",
+ Collections.emptyList(),
+ (namespace, name, bk) -> {
+ ViewEntity e =
createViewEntity(RandomIdGenerator.INSTANCE.nextId(), namespace, name);
+ bk.insert(e, false);
+ return new Object[] {e.nameIdentifier(), e.id()};
+ },
+ childIdent ->
+ ViewMetaService.getInstance()
+ .updateView(
+ childIdent,
+ entity -> {
+ ViewEntity v = (ViewEntity) entity;
+ return ViewEntity.builder()
+ .withId(v.id())
+ .withName(v.name())
+ .withNamespace(v.namespace())
+ .withAuditInfo(v.auditInfo())
+ .withColumns(v.columns())
+ .withRepresentations(v.representations())
+ .withComment("updated comment")
+ .build();
+ })),
+ new SchemaChildUpdateCase(
+ Entity.EntityType.FILESET,
+ "fileset_meta",
+ "fileset_id",
+ "fileset_version_info",
+ Collections.emptyList(),
+ (namespace, name, bk) -> {
+ FilesetEntity e =
+ createFilesetEntity(
+ RandomIdGenerator.INSTANCE.nextId(), namespace, name,
AUDIT_INFO);
+ bk.insert(e, false);
+ return new Object[] {e.nameIdentifier(), e.id()};
+ },
+ childIdent ->
+ FilesetMetaService.getInstance()
+ .updateFileset(
+ childIdent,
+ entity -> {
+ FilesetEntity f = (FilesetEntity) entity;
+ return FilesetEntity.builder()
+ .withId(f.id())
+ .withName(f.name())
+ .withNamespace(f.namespace())
+ .withFilesetType(f.filesetType())
+ .withStorageLocations(f.storageLocations())
+ .withAuditInfo(f.auditInfo())
+ .withComment("updated comment")
+ .withProperties(f.properties())
+ .build();
+ })),
+ new SchemaChildUpdateCase(
+ Entity.EntityType.FUNCTION,
+ "function_meta",
+ "function_id",
+ "function_version_info",
+ Collections.emptyList(),
+ (namespace, name, bk) -> {
+ FunctionEntity e =
+ createFunctionEntity(
+ RandomIdGenerator.INSTANCE.nextId(), namespace, name,
AUDIT_INFO);
+ bk.insert(e, false);
+ return new Object[] {e.nameIdentifier(), e.id()};
+ },
+ childIdent ->
+ FunctionMetaService.getInstance()
+ .updateFunction(
+ childIdent,
+ entity -> {
+ FunctionEntity fn = (FunctionEntity) entity;
+ return FunctionEntity.builder()
+ .withId(fn.id())
+ .withName(fn.name())
+ .withNamespace(fn.namespace())
+ .withAuditInfo(fn.auditInfo())
+ .withComment("updated comment")
+ .withFunctionType(fn.functionType())
+ .withDeterministic(fn.deterministic())
+ .withDefinitions(fn.definitions())
+ .build();
+ })),
+ new SchemaChildUpdateCase(
+ Entity.EntityType.MODEL,
+ "model_meta",
+ "model_id",
+ "model_version_info",
+ Collections.singletonList(new String[] {"model_version_alias_rel",
"model_id"}),
+ (namespace, name, bk) -> {
+ ModelEntity e =
+ createModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ name,
+ "model comment",
+ 0,
+ Collections.emptyMap(),
+ AUDIT_INFO);
+ bk.insert(e, false);
+ // Register version 0 with an alias so the version and alias
tables have rows
+ // that the cascade delete must clean up.
+ ModelVersionMetaService.getInstance()
+ .insertModelVersion(
+ ModelVersionEntity.builder()
+ .withModelIdentifier(e.nameIdentifier())
+ .withVersion(0)
+
.withUris(ImmutableMap.of(ModelVersion.URI_NAME_UNKNOWN, "/tmp"))
+ .withAliases(List.of("alias_" + name))
+ .withComment("version comment")
+ .withProperties(Collections.emptyMap())
+ .withAuditInfo(AUDIT_INFO)
+ .build());
+ return new Object[] {e.nameIdentifier(), e.id()};
+ },
+ childIdent ->
+ ModelMetaService.getInstance()
+ .updateModel(
+ childIdent,
+ entity -> {
+ ModelEntity m = (ModelEntity) entity;
+ return ModelEntity.builder()
+ .withId(m.id())
+ .withName(m.name())
+ .withNamespace(m.namespace())
+ .withAuditInfo(m.auditInfo())
+ .withComment("updated comment")
+ .withLatestVersion(m.latestVersion())
+ .withProperties(m.properties())
+ .build();
+ })),
+ new SchemaChildUpdateCase(
+ Entity.EntityType.TOPIC,
+ "topic_meta",
+ "topic_id",
+ null,
+ Collections.emptyList(),
+ (namespace, name, bk) -> {
+ TopicEntity e =
+ createTopicEntity(
+ RandomIdGenerator.INSTANCE.nextId(), namespace, name,
AUDIT_INFO);
+ bk.insert(e, false);
+ return new Object[] {e.nameIdentifier(), e.id()};
+ },
+ childIdent ->
+ TopicMetaService.getInstance()
+ .updateTopic(
+ childIdent,
+ entity -> {
+ TopicEntity topic = (TopicEntity) entity;
+ return TopicEntity.builder()
+ .withId(topic.id())
+ .withName(topic.name())
+ .withNamespace(topic.namespace())
+ .withAuditInfo(topic.auditInfo())
+ .withComment("updated comment")
+ .withProperties(topic.properties())
+ .build();
+ })));
+ }
+
+ private Map<Entity.EntityType, SchemaChildMove> schemaChildMoves() {
+ // Only table, view and function support cross-schema moves; those moves
take the
+ // catalog lock first.
+ Map<Entity.EntityType, SchemaChildMove> moves = new
EnumMap<>(Entity.EntityType.class);
+ moves.put(
+ Entity.EntityType.TABLE,
+ (childIdent, targetSchema) ->
+ TableMetaService.getInstance()
+ .updateTable(
+ childIdent,
+ entity -> {
+ TableEntity t = (TableEntity) entity;
+ return TableEntity.builder()
+ .withId(t.id())
+ .withName(t.name())
+ .withNamespace(Namespace.of(metalakeName,
catalogName, targetSchema))
+ .withAuditInfo(t.auditInfo())
+ .withColumns(t.columns())
+ .withComment("moved table comment")
+ .withProperties(t.properties())
+ .build();
+ }));
+ moves.put(
+ Entity.EntityType.VIEW,
+ (childIdent, targetSchema) ->
+ ViewMetaService.getInstance()
+ .updateView(
+ childIdent,
+ entity -> {
+ ViewEntity v = (ViewEntity) entity;
+ return ViewEntity.builder()
+ .withId(v.id())
+ .withName(v.name())
+ .withNamespace(Namespace.of(metalakeName,
catalogName, targetSchema))
+ .withAuditInfo(v.auditInfo())
+ .withColumns(v.columns())
+ .withRepresentations(v.representations())
+ .withComment(v.comment())
+ .build();
+ }));
+ moves.put(
+ Entity.EntityType.FUNCTION,
+ (childIdent, targetSchema) ->
+ FunctionMetaService.getInstance()
+ .updateFunction(
+ childIdent,
+ entity -> {
+ FunctionEntity fn = (FunctionEntity) entity;
+ return FunctionEntity.builder()
+ .withId(fn.id())
+ .withName(fn.name())
+ .withNamespace(Namespace.of(metalakeName,
catalogName, targetSchema))
+ .withAuditInfo(fn.auditInfo())
+ .withComment(fn.comment())
+ .withFunctionType(fn.functionType())
+ .withDeterministic(fn.deterministic())
+ .withDefinitions(fn.definitions())
+ .build();
+ }));
+ return moves;
+ }
+
+ @FunctionalInterface
+ private interface SchemaChildMove {
+ void run(NameIdentifier childIdentifier, String targetSchemaName) throws
Exception;
+ }
}