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 7f303cadad [#13002] fix(core): Fence owner assignment on the principal 
and keep one live owner per object (#13293)
7f303cadad is described below

commit 7f303cadad7623dfe8438ae8c2b9269f8e21e919
Author: Qi Yu <[email protected]>
AuthorDate: Mon Sep 21 21:53:29 2026 +0800

    [#13002] fix(core): Fence owner assignment on the principal and keep one 
live owner per object (#13293)
    
    ### What changes were proposed in this pull request?
    
    - Fence the metalake, owned object, and owner principal (User/Group) in
    one `setOwner` or `batchSetOwners` transaction. Assignments lock the
    owned object's row with `FOR UPDATE`, including when it has no owner
    yet. Batches lock objects in ID order. A metalake's own owner assignment
    takes an exclusive metalake lock.
    - Reject a principal that was deleted, replaced under the same name, or
    moved to another metalake while the assignment was waiting.
    - Keep the existing `owner_meta` unique key. The 2.0.0 upgrade scripts
    merge duplicate live owner rows left by older concurrent writes,
    retaining the row with the largest ID.
    
    ### Why are the changes needed?
    
    An assignment previously resolved the principal before its transaction
    and could insert a live owner row after another server deleted that
    principal. Two concurrent first assignments could also insert different
    live owners for the same object because no owner row existed to lock and
    the existing unique key includes `owner_id`. Locking the stable
    owned-object row serializes those assignments without adding a column or
    replacing the unique key.
    
    Fix: #13002
    
    ### Does this PR introduce _any_ user-facing change?
    
    No API change. An assignment waiting on a deleted object or principal
    now fails with the existing not-found error. Concurrent assignments
    through this service leave one live owner, with the later assignment
    taking effect.
    
    ### How was this patch tested?
    
    - `TestOwnerAssignmentWrites`: concurrent first assignments,
    reassignments, principal deletion and replacement, object deletion,
    metalake ownership, and batch rollback and retry.
    - `TestSQLScripts`: the 1.3.0 to 2.0.0 upgrade merges three live owners
    of one object while preserving historical rows that share a deletion
    timestamp.
    - `TestOwnerMetaService`, `TestOwnerAssignmentWrites`, and
    `TestSQLScripts` passed against H2, MySQL, and PostgreSQL.
    `:core:spotlessApply` passed.
---
 .../storage/relational/mapper/GroupMetaMapper.java |  12 +
 .../mapper/GroupMetaSQLProviderFactory.java        |   5 +
 .../storage/relational/mapper/OwnerMetaMapper.java |  16 +
 .../mapper/OwnerMetaSQLProviderFactory.java        |  82 ++++
 .../storage/relational/mapper/UserMetaMapper.java  |  12 +
 .../mapper/UserMetaSQLProviderFactory.java         |   5 +
 .../provider/base/GroupMetaBaseSQLProvider.java    |  12 +-
 .../provider/base/UserMetaBaseSQLProvider.java     |  12 +-
 .../mapper/provider/h2/GroupMetaH2Provider.java    |   6 +
 .../mapper/provider/h2/UserMetaH2Provider.java     |   6 +
 .../postgresql/GroupMetaPostgreSQLProvider.java    |   5 +
 .../postgresql/UserMetaPostgreSQLProvider.java     |   5 +
 .../relational/service/OwnerMetaService.java       | 147 +++++-
 .../apache/gravitino/storage/TestSQLScripts.java   |  75 +++
 .../service/TestOwnerAssignmentWrites.java         | 529 +++++++++++++++++++++
 scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql           |  13 +
 scripts/mysql/schema-2.0.0-mysql.sql               |   2 +-
 scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql     |  15 +
 .../upgrade-1.3.0-to-2.0.0-postgresql.sql          |  13 +
 19 files changed, 953 insertions(+), 19 deletions(-)

diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.java
index a5bd65d5c3..333be53597 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.java
@@ -20,6 +20,7 @@
 package org.apache.gravitino.storage.relational.mapper;
 
 import java.util.List;
+import javax.annotation.Nullable;
 import org.apache.gravitino.storage.relational.po.ExtendedGroupPO;
 import org.apache.gravitino.storage.relational.po.GroupPO;
 import org.apache.gravitino.storage.relational.po.auth.GroupUpdatedAt;
@@ -57,6 +58,17 @@ public interface GroupMetaMapper {
   @SelectProvider(type = GroupMetaSQLProviderFactory.class, method = 
"selectGroupMetaByIdForUpdate")
   GroupPO selectGroupMetaByIdForUpdate(@Param("groupId") Long groupId);
 
+  /**
+   * Returns an active group by ID and holds its lock for the current 
transaction.
+   *
+   * <p>The lock is shared on MySQL/PostgreSQL and exclusive on H2.
+   *
+   * @return the active group, or null if it does not exist
+   */
+  @Nullable
+  @SelectProvider(type = GroupMetaSQLProviderFactory.class, method = 
"selectGroupMetaByIdForShare")
+  GroupPO selectGroupMetaByIdForShare(@Param("groupId") Long groupId);
+
   @SelectProvider(
       type = GroupMetaSQLProviderFactory.class,
       method = "listExtendedGroupPOsByMetalakeIdAndNames")
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
index 14c6ab732e..215b4083c6 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
@@ -64,6 +64,11 @@ public class GroupMetaSQLProviderFactory {
     return getProvider().selectGroupMetaByIdForUpdate(groupId);
   }
 
+  /** Returns SQL that selects an active group by ID and locks it for shared 
access. */
+  public static String selectGroupMetaByIdForShare(@Param("groupId") Long 
groupId) {
+    return getProvider().selectGroupMetaByIdForShare(groupId);
+  }
+
   public static String listExtendedGroupPOsByMetalakeIdAndNames(
       @Param("metalakeId") Long metalakeId, @Param("groupNames") List<String> 
groupNames) {
     return getProvider().listExtendedGroupPOsByMetalakeIdAndNames(metalakeId, 
groupNames);
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
index e4ef4b6096..fc456f10cb 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
@@ -19,6 +19,8 @@
 package org.apache.gravitino.storage.relational.mapper;
 
 import java.util.List;
+import javax.annotation.Nullable;
+import org.apache.gravitino.Entity;
 import org.apache.gravitino.storage.relational.po.GroupOwnerRelPO;
 import org.apache.gravitino.storage.relational.po.GroupPO;
 import org.apache.gravitino.storage.relational.po.OwnerRelForDeletion;
@@ -44,6 +46,20 @@ public interface OwnerMetaMapper {
 
   String OWNER_TABLE_NAME = "owner_meta";
 
+  /**
+   * Locks an active metadata object before owner assignment.
+   *
+   * @return the object ID, or null if the object is no longer active
+   */
+  @Nullable
+  @SelectProvider(
+      type = OwnerMetaSQLProviderFactory.class,
+      method = "selectMetadataObjectIdForUpdate")
+  Long selectMetadataObjectIdForUpdate(
+      @Param("entityId") Long entityId,
+      @Param("metalakeId") Long metalakeId,
+      @Param("entityType") Entity.EntityType entityType);
+
   @SelectProvider(
       type = OwnerMetaSQLProviderFactory.class,
       method = "selectUserOwnerMetaByMetadataObjectIdAndType")
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
index 5957c6d91b..777904dabb 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
@@ -21,6 +21,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.Entity;
 import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType;
 import 
org.apache.gravitino.storage.relational.mapper.provider.base.OwnerMetaBaseSQLProvider;
 import 
org.apache.gravitino.storage.relational.mapper.provider.postgresql.OwnerMetaPostgreSQLProvider;
@@ -52,6 +53,87 @@ public class OwnerMetaSQLProviderFactory {
 
   static class OwnerMetaH2Provider extends OwnerMetaBaseSQLProvider {}
 
+  /** Returns SQL that locks an active metadata object before assigning its 
owner. */
+  public static String selectMetadataObjectIdForUpdate(
+      @Param("entityId") Long entityId,
+      @Param("metalakeId") Long metalakeId,
+      @Param("entityType") Entity.EntityType entityType) {
+    String table;
+    String idColumn;
+    switch (entityType) {
+      case CATALOG:
+        table = CatalogMetaMapper.TABLE_NAME;
+        idColumn = "catalog_id";
+        break;
+      case SCHEMA:
+        table = SchemaMetaMapper.TABLE_NAME;
+        idColumn = "schema_id";
+        break;
+      case TABLE:
+        table = TableMetaMapper.TABLE_NAME;
+        idColumn = "table_id";
+        break;
+      case COLUMN:
+        table = TableColumnMapper.COLUMN_TABLE_NAME;
+        idColumn = "column_id";
+        break;
+      case FILESET:
+        table = FilesetMetaMapper.META_TABLE_NAME;
+        idColumn = "fileset_id";
+        break;
+      case TOPIC:
+        table = TopicMetaMapper.TABLE_NAME;
+        idColumn = "topic_id";
+        break;
+      case MODEL:
+        table = ModelMetaMapper.TABLE_NAME;
+        idColumn = "model_id";
+        break;
+      case VIEW:
+        table = ViewMetaMapper.TABLE_NAME;
+        idColumn = "view_id";
+        break;
+      case FUNCTION:
+        table = FunctionMetaMapper.TABLE_NAME;
+        idColumn = "function_id";
+        break;
+      case ROLE:
+        table = RoleMetaMapper.ROLE_TABLE_NAME;
+        idColumn = "role_id";
+        break;
+      case TAG:
+        table = TagMetaMapper.TAG_TABLE_NAME;
+        idColumn = "tag_id";
+        break;
+      case POLICY:
+        table = PolicyMetaMapper.POLICY_META_TABLE_NAME;
+        idColumn = "policy_id";
+        break;
+      case JOB_TEMPLATE:
+        table = JobTemplateMetaMapper.TABLE_NAME;
+        idColumn = "job_template_id";
+        break;
+      case JOB:
+        table = JobMetaMapper.TABLE_NAME;
+        idColumn = "job_run_id";
+        break;
+      default:
+        throw new IllegalArgumentException("Unsupported owned object type: " + 
entityType);
+    }
+    // Column versions can share a column ID; always lock the same oldest live 
row.
+    String orderColumn = entityType == Entity.EntityType.COLUMN ? "id" : 
idColumn;
+    return "SELECT "
+        + idColumn
+        + " FROM "
+        + table
+        + " WHERE "
+        + idColumn
+        + " = #{entityId} AND metalake_id = #{metalakeId} AND deleted_at = 0"
+        + " ORDER BY "
+        + orderColumn
+        + " LIMIT 1 FOR UPDATE";
+  }
+
   public static String selectUserOwnerMetaByMetadataObjectIdAndType(
       @Param("metadataObjectId") Long metadataObjectId,
       @Param("metadataObjectType") String metadataObjectType) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.java
index ba79e5dca0..4a10ac8dca 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.java
@@ -20,6 +20,7 @@
 package org.apache.gravitino.storage.relational.mapper;
 
 import java.util.List;
+import javax.annotation.Nullable;
 import org.apache.gravitino.storage.relational.po.ExtendedUserPO;
 import org.apache.gravitino.storage.relational.po.UserPO;
 import org.apache.gravitino.storage.relational.po.auth.AuthPrefetchRow;
@@ -58,6 +59,17 @@ public interface UserMetaMapper {
   @SelectProvider(type = UserMetaSQLProviderFactory.class, method = 
"selectUserMetaByIdForUpdate")
   UserPO selectUserMetaByIdForUpdate(@Param("userId") Long userId);
 
+  /**
+   * Returns an active user by ID and holds its lock for the current 
transaction.
+   *
+   * <p>The lock is shared on MySQL/PostgreSQL and exclusive on H2.
+   *
+   * @return the active user, or null if it does not exist
+   */
+  @Nullable
+  @SelectProvider(type = UserMetaSQLProviderFactory.class, method = 
"selectUserMetaByIdForShare")
+  UserPO selectUserMetaByIdForShare(@Param("userId") Long userId);
+
   @InsertProvider(type = UserMetaSQLProviderFactory.class, method = 
"insertUserMeta")
   void insertUserMeta(@Param("userMeta") UserPO userPO);
 
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
index 0779d9920f..bc52e64d52 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
@@ -66,6 +66,11 @@ public class UserMetaSQLProviderFactory {
     return getProvider().selectUserMetaByIdForUpdate(userId);
   }
 
+  /** Returns SQL that selects an active user by ID and locks it for shared 
access. */
+  public static String selectUserMetaByIdForShare(@Param("userId") Long 
userId) {
+    return getProvider().selectUserMetaByIdForShare(userId);
+  }
+
   public static String insertUserMeta(@Param("userMeta") UserPO userPO) {
     return getProvider().insertUserMeta(userPO);
   }
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
index 002d68d2b5..4afa7ad3d2 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
@@ -141,13 +141,23 @@ public class GroupMetaBaseSQLProvider {
 
   /** Returns SQL that selects and locks an active group by ID. */
   public String selectGroupMetaByIdForUpdate(@Param("groupId") Long groupId) {
+    return selectGroupMetaById(groupId) + " FOR UPDATE";
+  }
+
+  /** Returns SQL that selects an active group by ID and locks it for shared 
access. */
+  public String selectGroupMetaByIdForShare(@Param("groupId") Long groupId) {
+    return selectGroupMetaById(groupId) + " LOCK IN SHARE MODE";
+  }
+
+  /** Returns SQL that selects an active group by ID. */
+  protected String selectGroupMetaById(Long groupId) {
     return "SELECT group_id as groupId, group_name as groupName,"
         + " metalake_id as metalakeId, audit_info as auditInfo,"
         + " current_version as currentVersion, last_version as lastVersion,"
         + " deleted_at as deletedAt"
         + " FROM "
         + GROUP_TABLE_NAME
-        + " WHERE group_id = #{groupId} AND deleted_at = 0 FOR UPDATE";
+        + " WHERE group_id = #{groupId} AND deleted_at = 0";
   }
 
   public String listExtendedGroupPOsByMetalakeIdAndNames(
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
index 98e8931f1d..dbe362b7db 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
@@ -55,13 +55,23 @@ public class UserMetaBaseSQLProvider {
 
   /** Returns SQL that selects and locks an active user by ID. */
   public String selectUserMetaByIdForUpdate(@Param("userId") Long userId) {
+    return selectUserMetaById(userId) + " FOR UPDATE";
+  }
+
+  /** Returns SQL that selects an active user by ID and locks it for shared 
access. */
+  public String selectUserMetaByIdForShare(@Param("userId") Long userId) {
+    return selectUserMetaById(userId) + " LOCK IN SHARE MODE";
+  }
+
+  /** Returns SQL that selects an active user by ID. */
+  protected String selectUserMetaById(Long userId) {
     return "SELECT user_id as userId, user_name as userName,"
         + " metalake_id as metalakeId,"
         + " audit_info as auditInfo, current_version as currentVersion,"
         + " last_version as lastVersion, deleted_at as deletedAt"
         + " FROM "
         + USER_TABLE_NAME
-        + " WHERE user_id = #{userId} AND deleted_at = 0 FOR UPDATE";
+        + " WHERE user_id = #{userId} AND deleted_at = 0";
   }
 
   public String insertUserMeta(@Param("userMeta") UserPO userPO) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
index 2d542428cd..e25aeb4592 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
@@ -28,6 +28,12 @@ import 
org.apache.gravitino.storage.relational.mapper.provider.base.GroupMetaBas
 import org.apache.ibatis.annotations.Param;
 
 public class GroupMetaH2Provider extends GroupMetaBaseSQLProvider {
+  @Override
+  public String selectGroupMetaByIdForShare(Long groupId) {
+    // H2 has no shared row-lock syntax, matching the other parent-fencing 
providers.
+    return selectGroupMetaByIdForUpdate(groupId);
+  }
+
   @Override
   public String listExtendedGroupPOsByMetalakeId(@Param("metalakeId") Long 
metalakeId) {
     return "SELECT gt.group_id as groupId, gt.group_name as groupName,"
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
index 83fe6774b5..9e28a2ad9f 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
@@ -27,6 +27,12 @@ import 
org.apache.gravitino.storage.relational.mapper.provider.base.UserMetaBase
 import org.apache.ibatis.annotations.Param;
 
 public class UserMetaH2Provider extends UserMetaBaseSQLProvider {
+  @Override
+  public String selectUserMetaByIdForShare(Long userId) {
+    // H2 has no shared row-lock syntax, matching the other parent-fencing 
providers.
+    return selectUserMetaByIdForUpdate(userId);
+  }
+
   @Override
   public String listExtendedUserPOsByMetalakeId(@Param("metalakeId") Long 
metalakeId) {
     return "SELECT ut.user_id as userId, ut.user_name as userName,"
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
index 1a24acec0f..bac95d6f99 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
@@ -30,6 +30,11 @@ import org.apache.gravitino.storage.relational.po.GroupPO;
 import org.apache.ibatis.annotations.Param;
 
 public class GroupMetaPostgreSQLProvider extends GroupMetaBaseSQLProvider {
+  @Override
+  public String selectGroupMetaByIdForShare(Long groupId) {
+    return selectGroupMetaById(groupId) + " FOR SHARE";
+  }
+
   @Override
   public String softDeleteGroupMetaByGroupId(Long groupId, Long 
currentVersion) {
     return "UPDATE "
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
index 4562901e8b..5ec9ef7030 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
@@ -29,6 +29,11 @@ import org.apache.gravitino.storage.relational.po.UserPO;
 import org.apache.ibatis.annotations.Param;
 
 public class UserMetaPostgreSQLProvider extends UserMetaBaseSQLProvider {
+  @Override
+  public String selectUserMetaByIdForShare(Long userId) {
+    return selectUserMetaById(userId) + " FOR SHARE";
+  }
+
   @Override
   public String softDeleteUserMetaByUserId(Long userId, Long currentVersion) {
     return "UPDATE "
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
index 2b261e35b3..4221254694 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
@@ -23,6 +23,7 @@ import static 
org.apache.gravitino.metrics.source.MetricsSource.GRAVITINO_RELATI
 import com.google.common.base.Preconditions;
 import java.util.ArrayList;
 import java.util.Collections;
+import java.util.Comparator;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -38,7 +39,10 @@ import org.apache.gravitino.authorization.AuthorizationUtils;
 import org.apache.gravitino.meta.GroupEntity;
 import org.apache.gravitino.meta.UserEntity;
 import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.mapper.GroupMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.MetalakeMetaMapper;
 import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.UserMetaMapper;
 import org.apache.gravitino.storage.relational.po.GroupOwnerRelPO;
 import org.apache.gravitino.storage.relational.po.GroupPO;
 import org.apache.gravitino.storage.relational.po.OwnerRelForDeletion;
@@ -52,10 +56,10 @@ import org.apache.gravitino.utils.NameIdentifierUtil;
 /** This class is an utilization class to retrieve owner relation. */
 public class OwnerMetaService {
 
-  private OwnerMetaService() {}
-
   private static final OwnerMetaService INSTANCE = new OwnerMetaService();
 
+  private OwnerMetaService() {}
+
   public static OwnerMetaService getInstance() {
     return INSTANCE;
   }
@@ -173,9 +177,8 @@ public class OwnerMetaService {
       Entity.EntityType entityType,
       NameIdentifier owner,
       Entity.EntityType ownerType) {
-    long metalakeId =
-        MetalakeMetaService.getInstance()
-            .getMetalakeIdByName(NameIdentifierUtil.getMetalake(entity));
+    String metalake = NameIdentifierUtil.getMetalake(entity);
+    long metalakeId = 
MetalakeMetaService.getInstance().getMetalakeIdByName(metalake);
 
     Long entityId = EntityIdService.getEntityId(entity, entityType);
     Long ownerId = EntityIdService.getEntityId(owner, ownerType);
@@ -183,14 +186,22 @@ public class OwnerMetaService {
     OwnerRelPO ownerRelPO =
         POConverters.initializeOwnerRelPOsWithVersion(
             metalakeId, ownerType.name(), ownerId, entityType.name(), 
entityId);
-    SessionUtils.doMultipleWithCommit(
+    String metadataObjectType =
+        NameIdentifierUtil.toMetadataObject(entity, entityType).type().name();
+    assignOwner(
+        metalake,
+        metalakeId,
+        owner,
+        ownerType,
+        ownerId,
+        List.of(entityId),
+        entityType,
         () ->
             SessionUtils.doWithoutCommit(
                 OwnerMetaMapper.class,
                 mapper ->
                     mapper.softDeleteOwnerRelByMetadataObjectIdAndType(
-                        entityId,
-                        NameIdentifierUtil.toMetadataObject(entity, 
entityType).type().name())),
+                        entityId, metadataObjectType)),
         () ->
             SessionUtils.doWithoutCommit(
                 OwnerMetaMapper.class, mapper -> 
mapper.insertOwnerRel(ownerRelPO)));
@@ -218,20 +229,33 @@ public class OwnerMetaService {
     long metalakeId = 
MetalakeMetaService.getInstance().getMetalakeIdByName(metalake);
     Long ownerId = EntityIdService.getEntityId(ownerIdent, ownerType);
 
-    List<OwnerRelForDeletion> deletions = new ArrayList<>(ownedObjects.size());
-    List<OwnerRelPO> ownerRelPOs = new ArrayList<>(ownedObjects.size());
+    // Resolve every object first and write in stable id order, so two batches 
that overlap take
+    // the same row locks in the same order and cannot deadlock each other.
+    List<Long> entityIds = new ArrayList<>(ownedObjects.size());
     for (NameIdentifier entity : ownedObjects) {
-      Long entityId = EntityIdService.getEntityId(entity, ownedObjectType);
-      deletions.add(
-          new OwnerRelForDeletion(
-              entityId,
-              NameIdentifierUtil.toMetadataObject(entity, 
ownedObjectType).type().name()));
+      entityIds.add(EntityIdService.getEntityId(entity, ownedObjectType));
+    }
+    entityIds.sort(Comparator.naturalOrder());
+    String metadataObjectType =
+        NameIdentifierUtil.toMetadataObject(ownedObjects.get(0), 
ownedObjectType).type().name();
+
+    List<OwnerRelForDeletion> deletions = new ArrayList<>(entityIds.size());
+    List<OwnerRelPO> ownerRelPOs = new ArrayList<>(entityIds.size());
+    for (Long entityId : entityIds) {
+      deletions.add(new OwnerRelForDeletion(entityId, metadataObjectType));
       ownerRelPOs.add(
           POConverters.initializeOwnerRelPOsWithVersion(
               metalakeId, ownerType.name(), ownerId, ownedObjectType.name(), 
entityId));
     }
 
-    SessionUtils.doMultipleWithCommit(
+    assignOwner(
+        metalake,
+        metalakeId,
+        ownerIdent,
+        ownerType,
+        ownerId,
+        entityIds,
+        ownedObjectType,
         () ->
             SessionUtils.doWithoutCommit(
                 OwnerMetaMapper.class,
@@ -240,4 +264,95 @@ public class OwnerMetaService {
             SessionUtils.doWithoutCommit(
                 OwnerMetaMapper.class, mapper -> 
mapper.batchInsertOwnerRels(ownerRelPOs)));
   }
+
+  /**
+   * Serializes assignments on the owned object's row, including the first 
assignment when no owner
+   * relation exists. The metalake is locked first to fence cascade deletion; 
object rows are locked
+   * in stable ID order; the principal is then fenced before owner relations 
change.
+   */
+  private void assignOwner(
+      String metalake,
+      long metalakeId,
+      NameIdentifier owner,
+      Entity.EntityType ownerType,
+      long ownerId,
+      List<Long> entityIds,
+      Entity.EntityType ownedObjectType,
+      Runnable retirePreviousOwners,
+      Runnable insertOwners) {
+    SessionUtils.doMultipleWithCommit(
+        () ->
+            lockMetalakeForOwnerWrite(
+                metalake, metalakeId, ownedObjectType == 
Entity.EntityType.METALAKE),
+        () -> lockOwnedObjectsForOwnerWrite(entityIds, ownedObjectType, 
metalakeId),
+        () -> lockPrincipalForOwnerWrite(owner, ownerType, ownerId, 
metalakeId),
+        retirePreviousOwners,
+        insertOwners);
+  }
+
+  private void lockMetalakeForOwnerWrite(String metalake, long metalakeId, 
boolean exclusive) {
+    OccWriteSupport.lockParentForChildWrite(
+        metalake,
+        Entity.EntityType.METALAKE,
+        () ->
+            SessionUtils.getWithoutCommit(
+                MetalakeMetaMapper.class,
+                mapper ->
+                    exclusive
+                        ? mapper.selectMetalakeMetaByIdForUpdate(metalakeId)
+                        : mapper.selectMetalakeMetaByIdForShare(metalakeId)),
+        null,
+        current -> Objects.equals(current.getMetalakeName(), metalake));
+  }
+
+  private void lockOwnedObjectsForOwnerWrite(
+      List<Long> entityIds, Entity.EntityType entityType, long metalakeId) {
+    if (entityType == Entity.EntityType.METALAKE) {
+      return;
+    }
+    for (Long entityId : entityIds) {
+      OccWriteSupport.lockParentForChildWrite(
+          String.valueOf(entityId),
+          entityType,
+          () ->
+              SessionUtils.getWithoutCommit(
+                  OwnerMetaMapper.class,
+                  mapper ->
+                      mapper.selectMetadataObjectIdForUpdate(entityId, 
metalakeId, entityType)),
+          null,
+          current -> Objects.equals(current, entityId));
+    }
+  }
+
+  private void lockPrincipalForOwnerWrite(
+      NameIdentifier owner, Entity.EntityType ownerType, long ownerId, long 
metalakeId) {
+    switch (ownerType) {
+      case USER:
+        OccWriteSupport.lockParentForChildWrite(
+            owner.name(),
+            ownerType,
+            () ->
+                SessionUtils.getWithoutCommit(
+                    UserMetaMapper.class, mapper -> 
mapper.selectUserMetaByIdForShare(ownerId)),
+            null,
+            current ->
+                Objects.equals(current.getMetalakeId(), metalakeId)
+                    && Objects.equals(current.getUserName(), owner.name()));
+        return;
+      case GROUP:
+        OccWriteSupport.lockParentForChildWrite(
+            owner.name(),
+            ownerType,
+            () ->
+                SessionUtils.getWithoutCommit(
+                    GroupMetaMapper.class, mapper -> 
mapper.selectGroupMetaByIdForShare(ownerId)),
+            null,
+            current ->
+                Objects.equals(current.getMetalakeId(), metalakeId)
+                    && Objects.equals(current.getGroupName(), owner.name()));
+        return;
+      default:
+        throw new IllegalArgumentException("Unsupported owner type: " + 
ownerType);
+    }
+  }
 }
diff --git 
a/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java 
b/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
index ba06be927c..c758614778 100644
--- a/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
+++ b/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
@@ -126,6 +126,81 @@ public class TestSQLScripts extends TestJDBCBackend {
     }
   }
 
+  /**
+   * The owner unique key allows one live row per (owner, object). Rows left 
by concurrent
+   * assignments are merged during the upgrade: the newest live row (largest 
id) stays, while older
+   * ones are soft-deleted.
+   */
+  @TestTemplate
+  public void testUpgradeToTwoZeroMergesDuplicateLiveOwners() throws 
SQLException, IOException {
+    String gravitinoHome = System.getenv("GRAVITINO_HOME");
+    Assertions.assertNotNull(gravitinoHome, "GRAVITINO_HOME environment 
variable is not set");
+    Path scriptDir = Path.of(gravitinoHome, "scripts", 
backendType.toLowerCase());
+    String suffix = "-" + backendType.toLowerCase() + ".sql";
+    dropAllTables();
+    executeScript(scriptDir.resolve("schema-1.3.0" + suffix).toFile());
+
+    String insert =
+        "INSERT INTO owner_meta (id, metalake_id, owner_id, owner_type, 
metadata_object_id,"
+            + " metadata_object_type, audit_info, current_version, 
last_version, deleted_at,"
+            + " updated_at) VALUES (%d, 1, %d, 'USER', %d, 'CATALOG', '{}', 1, 
1, %d, 0)";
+    List<String> rows =
+        List.of(
+            // Three owners of object 10 leave two rows to retire in the same 
statement.
+            String.format(insert, 1, 100, 10, 0),
+            String.format(insert, 2, 200, 10, 0),
+            String.format(insert, 3, 300, 10, 0),
+            // Historical rows may already share a deletion timestamp.
+            String.format(insert, 4, 100, 20, 0),
+            String.format(insert, 5, 200, 20, 5),
+            String.format(insert, 6, 300, 20, 5),
+            // Object 10 as a SCHEMA is a different object.
+            String.format(insert, 7, 400, 10, 0).replace("'CATALOG'", 
"'SCHEMA'"));
+    try (SqlSession sqlSession =
+            
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+        Connection connection = sqlSession.getConnection();
+        Statement statement = connection.createStatement()) {
+      for (String row : rows) {
+        statement.execute(row);
+      }
+    }
+
+    executeScript(scriptDir.resolve("upgrade-1.3.0-to-2.0.0" + 
suffix).toFile());
+
+    try (SqlSession sqlSession =
+            
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+        Connection connection = sqlSession.getConnection();
+        Statement statement = connection.createStatement();
+        ResultSet live =
+            statement.executeQuery("SELECT id FROM owner_meta WHERE deleted_at 
= 0 ORDER BY id")) {
+      List<Long> liveIds = new ArrayList<>();
+      while (live.next()) {
+        liveIds.add(live.getLong(1));
+      }
+      Assertions.assertEquals(List.of(3L, 4L, 7L), liveIds);
+    }
+    try (SqlSession sqlSession =
+            
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+        Connection connection = sqlSession.getConnection();
+        Statement statement = connection.createStatement();
+        ResultSet retired =
+            statement.executeQuery("SELECT deleted_at, updated_at FROM 
owner_meta WHERE id = 1")) {
+      Assertions.assertTrue(retired.next());
+      Assertions.assertTrue(retired.getLong(1) > 0, "older duplicate must be 
soft-deleted");
+      Assertions.assertEquals(retired.getLong(1), retired.getLong(2));
+    }
+    try (SqlSession sqlSession =
+            
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+        Connection connection = sqlSession.getConnection();
+        Statement statement = connection.createStatement()) {
+      // Historical rows can share a deletion timestamp after the upgrade.
+      statement.execute(String.format(insert, 8, 500, 20, 5));
+      // The existing key still rejects a second live row for the same owner 
and object.
+      Assertions.assertThrows(
+          SQLException.class, () -> statement.execute(String.format(insert, 9, 
100, 20, 0)));
+    }
+  }
+
   private void executeScript(File scriptFile) throws IOException, SQLException 
{
     List<String> ddls = extractStatements(scriptFile.toPath());
     try (SqlSession sqlSession =
diff --git 
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOwnerAssignmentWrites.java
 
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOwnerAssignmentWrites.java
new file mode 100644
index 0000000000..9c0e92dc75
--- /dev/null
+++ 
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOwnerAssignmentWrites.java
@@ -0,0 +1,529 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *  http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.gravitino.storage.relational.service;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+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.util.ArrayList;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import javax.annotation.Nullable;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.authorization.AuthorizationUtils;
+import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.meta.CatalogEntity;
+import org.apache.gravitino.meta.GroupEntity;
+import org.apache.gravitino.meta.UserEntity;
+import org.apache.gravitino.storage.RandomIdGenerator;
+import org.apache.gravitino.storage.relational.TestJDBCBackend;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.gravitino.storage.relational.session.SqlSessions;
+import org.apache.gravitino.storage.relational.utils.SessionUtils;
+import org.apache.ibatis.session.SqlSession;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.function.Executable;
+
+/**
+ * Races between owner assignment and deletion of the owned object or owner 
principal, and between
+ * two assignments on the same object. Every scenario is driven by a real 
transaction held open on
+ * one thread while the contender runs on another, so the assertions describe 
what two Gravitino
+ * servers sharing one database would observe.
+ */
+class TestOwnerAssignmentWrites extends TestJDBCBackend {
+  private static final String METALAKE = "owner_write_metalake";
+
+  @TestTemplate
+  void testAssignmentWaitsForUncommittedPrincipalDeleteAndFails() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    for (boolean group : List.of(false, true)) {
+      CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned_" + group);
+      NameIdentifier principal = insertPrincipal(group, "deleted_" + group);
+      long id = principalId(group, principal);
+      Throwable failure =
+          whileTransactionHeld(
+              () -> deletePrincipal(group, principal),
+              () -> setOwner(owned, principal, type(group)));
+      assertInstanceOf(NoSuchEntityException.class, failure);
+      assertFalse(owner(owned).isPresent());
+      assertEquals(0, liveOwnerRows(id));
+    }
+  }
+
+  @TestTemplate
+  void testAssignmentWaitsForUncommittedOwnedObjectDeleteAndFails() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned");
+    NameIdentifier principal = insertPrincipal(false, "owner");
+    Throwable failure =
+        whileTransactionHeld(
+            () ->
+                assertTrue(
+                    
CatalogMetaService.getInstance().deleteCatalog(owned.nameIdentifier(), false)),
+            () -> setOwner(owned, principal, Entity.EntityType.USER));
+    assertInstanceOf(NoSuchEntityException.class, failure);
+    assertEquals(0, liveOwnerRowsForObject(owned));
+  }
+
+  @TestTemplate
+  void testPrincipalDeleteWaitsForUncommittedAssignmentAndCleansIt() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    for (boolean group : List.of(false, true)) {
+      CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned_" + group);
+      NameIdentifier principal = insertPrincipal(group, "deleted_" + group);
+      long id = principalId(group, principal);
+      assertNull(
+          whileTransactionHeld(
+              () -> setOwner(owned, principal, type(group)),
+              () -> deletePrincipal(group, principal)));
+      assertFalse(owner(owned).isPresent());
+      assertEquals(0, liveOwnerRows(id));
+    }
+  }
+
+  @TestTemplate
+  void testAssignmentRejectsSameNameReplacementObservedByStaleId() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    for (boolean group : List.of(false, true)) {
+      CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned_" + group);
+      NameIdentifier principal = insertPrincipal(group, "recreated_" + group);
+      long staleId = principalId(group, principal);
+      Throwable failure =
+          whileTransactionHeld(
+              () -> {
+                deletePrincipal(group, principal);
+                insertPrincipal(group, principal.name());
+              },
+              () -> setOwner(owned, principal, type(group)));
+      assertInstanceOf(NoSuchEntityException.class, failure);
+      assertEquals(0, liveOwnerRows(staleId));
+      assertFalse(owner(owned).isPresent());
+
+      // A fresh lookup resolves the replacement and succeeds.
+      setOwner(owned, principal, type(group));
+      assertEquals(principal.name(), ownerName(owned));
+      assertEquals(1, liveOwnerRows(principalId(group, principal)));
+    }
+  }
+
+  @TestTemplate
+  void testConcurrentInitialAssignmentsLeaveExactlyOneLiveOwner() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned");
+    NameIdentifier first = insertPrincipal(false, "first");
+    NameIdentifier second = insertPrincipal(true, "second");
+    assertNull(
+        whileTransactionHeld(
+            () -> setOwner(owned, first, Entity.EntityType.USER),
+            () -> setOwner(owned, second, Entity.EntityType.GROUP),
+            true));
+    assertEquals(1, liveOwnerRowsForObject(owned));
+    assertEquals("second", ownerName(owned));
+  }
+
+  @TestTemplate
+  void testConcurrentMetalakeAssignmentsSerializeOnTheMetalakeRow() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    NameIdentifier owned = NameIdentifier.of(METALAKE);
+    NameIdentifier first = insertPrincipal(false, "first");
+    NameIdentifier second = insertPrincipal(true, "second");
+    assertNull(
+        whileTransactionHeld(
+            () ->
+                OwnerMetaService.getInstance()
+                    .setOwner(owned, Entity.EntityType.METALAKE, first, 
Entity.EntityType.USER),
+            () ->
+                OwnerMetaService.getInstance()
+                    .setOwner(owned, Entity.EntityType.METALAKE, second, 
Entity.EntityType.GROUP)));
+    long metalakeId = 
MetalakeMetaService.getInstance().getMetalakeIdByName(METALAKE);
+    assertEquals(
+        1,
+        queryLong(
+            "SELECT COUNT(*) FROM owner_meta WHERE metadata_object_id = "
+                + metalakeId
+                + " AND metadata_object_type = 'METALAKE' AND deleted_at = 
0"));
+    assertEquals(
+        "second",
+        assertInstanceOf(
+                GroupEntity.class,
+                OwnerMetaService.getInstance()
+                    .getOwner(owned, Entity.EntityType.METALAKE)
+                    .orElseThrow())
+            .name());
+  }
+
+  @TestTemplate
+  void testConcurrentReassignmentsSerializeOnTheExistingOwnerRow() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned");
+    NameIdentifier initial = insertPrincipal(false, "initial");
+    NameIdentifier first = insertPrincipal(false, "first");
+    NameIdentifier second = insertPrincipal(true, "second");
+    setOwner(owned, initial, Entity.EntityType.USER);
+    assertNull(
+        whileTransactionHeld(
+            () -> setOwner(owned, first, Entity.EntityType.USER),
+            () -> setOwner(owned, second, Entity.EntityType.GROUP),
+            true));
+    assertEquals(1, liveOwnerRowsForObject(owned));
+    assertEquals("second", ownerName(owned));
+  }
+
+  @TestTemplate
+  void testBatchAssignmentRollsBackWhenPrincipalIsDeleted() throws Exception {
+    createAndInsertMakeLake(METALAKE);
+    for (boolean group : List.of(false, true)) {
+      NameIdentifier previous = insertPrincipal(false, "previous_" + group);
+      NameIdentifier principal = insertPrincipal(group, "deleted_" + group);
+      long id = principalId(group, principal);
+      List<CatalogEntity> owned = new ArrayList<>();
+      for (int i = 0; i < 3; i++) {
+        CatalogEntity catalog = createAndInsertCatalog(METALAKE, "owned_" + 
group + "_" + i);
+        owned.add(catalog);
+        if (i > 0) {
+          setOwner(catalog, previous, Entity.EntityType.USER);
+        }
+      }
+      Throwable failure =
+          whileTransactionHeld(
+              () -> deletePrincipal(group, principal),
+              () -> batchSetOwners(owned, principal, type(group)));
+      assertInstanceOf(NoSuchEntityException.class, failure);
+      assertFalse(SessionUtils.isInTransaction());
+      assertFalse(owner(owned.get(0)).isPresent());
+      assertEquals(previous.name(), ownerName(owned.get(1)));
+      assertEquals(previous.name(), ownerName(owned.get(2)));
+      assertEquals(0, liveOwnerRows(id));
+    }
+  }
+
+  @TestTemplate
+  void testBatchAssignmentSucceedsAfterPrincipalDeleteRollsBack() throws 
Exception {
+    createAndInsertMakeLake(METALAKE);
+    for (boolean group : List.of(false, true)) {
+      NameIdentifier principal = insertPrincipal(group, "kept_" + group);
+      List<CatalogEntity> owned =
+          List.of(
+              createAndInsertCatalog(METALAKE, "owned_" + group + "_0"),
+              createAndInsertCatalog(METALAKE, "owned_" + group + "_1"));
+      assertNull(
+          whileTransactionHeld(
+              () -> deletePrincipal(group, principal),
+              () -> batchSetOwners(owned, principal, type(group)),
+              () -> {},
+              false));
+      for (CatalogEntity catalog : owned) {
+        assertEquals(principal.name(), ownerName(catalog));
+      }
+      assertEquals(2, liveOwnerRows(principalId(group, principal)));
+    }
+  }
+
+  private void setOwner(CatalogEntity owned, NameIdentifier owner, 
Entity.EntityType ownerType) {
+    OwnerMetaService.getInstance()
+        .setOwner(owned.nameIdentifier(), Entity.EntityType.CATALOG, owner, 
ownerType);
+  }
+
+  private void batchSetOwners(
+      List<CatalogEntity> owned, NameIdentifier owner, Entity.EntityType 
ownerType) {
+    List<NameIdentifier> identifiers = new ArrayList<>();
+    for (CatalogEntity catalog : owned) {
+      identifiers.add(catalog.nameIdentifier());
+    }
+    OwnerMetaService.getInstance()
+        .batchSetOwners(identifiers, Entity.EntityType.CATALOG, owner, 
ownerType);
+  }
+
+  private Optional<Entity> owner(CatalogEntity owned) {
+    return OwnerMetaService.getInstance()
+        .getOwner(owned.nameIdentifier(), Entity.EntityType.CATALOG);
+  }
+
+  private String ownerName(CatalogEntity owned) {
+    Entity entity = owner(owned).orElseThrow(() -> new AssertionError("No 
owner"));
+    return entity instanceof UserEntity
+        ? ((UserEntity) entity).name()
+        : ((GroupEntity) entity).name();
+  }
+
+  private NameIdentifier insertPrincipal(boolean group, String name) throws 
IOException {
+    long id = RandomIdGenerator.INSTANCE.nextId();
+    if (group) {
+      GroupMetaService.getInstance()
+          .insertGroup(
+              createGroupEntity(
+                  id, AuthorizationUtils.ofGroupNamespace(METALAKE), name, 
AUDIT_INFO, null, null),
+              false);
+      return AuthorizationUtils.ofGroup(METALAKE, name);
+    }
+    UserMetaService.getInstance()
+        .insertUser(
+            createUserEntity(id, AuthorizationUtils.ofUserNamespace(METALAKE), 
name, AUDIT_INFO),
+            false);
+    return AuthorizationUtils.ofUser(METALAKE, name);
+  }
+
+  private void deletePrincipal(boolean group, NameIdentifier principal) {
+    if (group) {
+      assertTrue(GroupMetaService.getInstance().deleteGroup(principal));
+    } else {
+      assertTrue(UserMetaService.getInstance().deleteUser(principal));
+    }
+  }
+
+  private long principalId(boolean group, NameIdentifier principal) {
+    return EntityIdService.getEntityId(principal, type(group));
+  }
+
+  private Entity.EntityType type(boolean group) {
+    return group ? Entity.EntityType.GROUP : Entity.EntityType.USER;
+  }
+
+  private long liveOwnerRows(long ownerId) throws Exception {
+    return queryLong(
+        "SELECT COUNT(*) FROM owner_meta WHERE owner_id = " + ownerId + " AND 
deleted_at = 0");
+  }
+
+  private long liveOwnerRowsForObject(CatalogEntity owned) throws Exception {
+    return queryLong(
+        "SELECT COUNT(*) FROM owner_meta WHERE metadata_object_id = "
+            + owned.id()
+            + " AND metadata_object_type = 'CATALOG' AND deleted_at = 0");
+  }
+
+  private long queryLong(String sql) throws Exception {
+    try (SqlSession session =
+            
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+        Statement statement = session.getConnection().createStatement();
+        ResultSet rows = statement.executeQuery(sql)) {
+      assertTrue(rows.next());
+      return rows.getLong(1);
+    }
+  }
+
+  private Throwable whileTransactionHeld(Executable holder, Executable 
contender) throws Exception {
+    return whileTransactionHeld(holder, contender, () -> {}, true, false);
+  }
+
+  private Throwable whileTransactionHeld(
+      Executable holder, Executable contender, boolean standaloneContender) 
throws Exception {
+    return whileTransactionHeld(holder, contender, () -> {}, true, 
standaloneContender);
+  }
+
+  private Throwable whileTransactionHeld(
+      Executable holder, Executable contender, Executable beforeCompletion, 
boolean commitHolder)
+      throws Exception {
+    return whileTransactionHeld(holder, contender, beforeCompletion, 
commitHolder, false);
+  }
+
+  /**
+   * Runs {@code holder} inside a transaction held open on this thread, starts 
{@code contender} on
+   * another thread, waits until the database reports the contender blocked on 
the holder, then
+   * commits or rolls back the holder and returns the contender's failure, or 
null when it
+   * committed.
+   *
+   * <p>By default the contender is wrapped in a transaction of its own so 
that its session can be
+   * identified. A {@code standaloneContender} runs exactly as production 
does, owning its
+   * transactions, which is what the owner assignment needs to replay a lost 
race: its session is
+   * then not known in advance and the wait is recognised by the holder side 
alone.
+   */
+  private Throwable whileTransactionHeld(
+      Executable holder,
+      Executable contender,
+      Executable beforeCompletion,
+      boolean commitHolder,
+      boolean standaloneContender)
+      throws Exception {
+    ExecutorService executor = Executors.newSingleThreadExecutor();
+    CompletableFuture<Long> started = new CompletableFuture<>();
+    if (standaloneContender) {
+      raiseDefaultLockTimeout();
+    }
+    SessionUtils.beginTransaction();
+    try {
+      long holderId = prepareTransaction();
+      Assertions.assertDoesNotThrow(holder);
+      Future<Throwable> result =
+          standaloneContender
+              ? executor.submit(
+                  () -> {
+                    try {
+                      contender.execute();
+                      return null;
+                    } catch (Throwable failure) {
+                      return failure;
+                    }
+                  })
+              : submitTransaction(executor, started, contender);
+      awaitBlockedBy(
+          result, standaloneContender ? null : started.get(10, 
TimeUnit.SECONDS), holderId);
+      Assertions.assertDoesNotThrow(beforeCompletion);
+      if (commitHolder) {
+        SessionUtils.commitTransaction();
+      } else {
+        SessionUtils.rollbackTransaction();
+      }
+      return result.get(10, TimeUnit.SECONDS);
+    } finally {
+      SessionUtils.rollbackTransaction();
+      executor.shutdownNow();
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+    }
+  }
+
+  private Future<Throwable> submitTransaction(
+      ExecutorService executor, CompletableFuture<Long> started, Executable 
operation) {
+    return executor.submit(
+        () -> {
+          SessionUtils.beginTransaction();
+          try {
+            started.complete(prepareTransaction());
+            operation.execute();
+            SessionUtils.commitTransaction();
+            return null;
+          } catch (Throwable failure) {
+            started.completeExceptionally(failure);
+            return failure;
+          } finally {
+            SessionUtils.rollbackTransaction();
+          }
+        });
+  }
+
+  private long prepareTransaction() throws SQLException {
+    SqlSession session = SqlSessions.getSqlSession();
+    try (Statement statement = session.getConnection().createStatement()) {
+      String sessionIdQuery;
+      switch (backendType) {
+        case "h2":
+          // Keep the engine timeout above the test's lock-observation 
deadline.
+          statement.execute("SET LOCK_TIMEOUT 30000");
+          sessionIdQuery = "SELECT SESSION_ID()";
+          break;
+        case "mysql":
+          statement.execute("SET SESSION innodb_lock_wait_timeout = 30");
+          sessionIdQuery = "SELECT CONNECTION_ID()";
+          break;
+        case "postgresql":
+          statement.execute("SET LOCAL lock_timeout = '30s'");
+          sessionIdQuery = "SELECT pg_backend_pid()";
+          break;
+        default:
+          throw new IllegalStateException("Unsupported backend: " + 
backendType);
+      }
+      try (ResultSet rows = statement.executeQuery(sessionIdQuery)) {
+        assertTrue(rows.next());
+        return rows.getLong(1);
+      }
+    } finally {
+      SqlSessions.closeSqlSession();
+    }
+  }
+
+  /**
+   * H2 gives new sessions a one-second lock timeout. A standalone contender 
opens its own sessions,
+   * so the database-wide default is raised instead of a per-session setting.
+   */
+  private void raiseDefaultLockTimeout() throws SQLException {
+    if (!"h2".equals(backendType)) {
+      return;
+    }
+    try (SqlSession session =
+            
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+        Statement statement = session.getConnection().createStatement()) {
+      statement.execute("SET DEFAULT_LOCK_TIMEOUT 30000");
+    }
+  }
+
+  /**
+   * Waits until the database reports a session blocked by the holder: the 
given contender session
+   * when known, otherwise any session. On MySQL a waiter on the implicit lock 
of an uncommitted
+   * insert is reported as blocked by itself rather than by the inserting 
transaction, so that shape
+   * is accepted too.
+   */
+  private void awaitBlockedBy(Future<Throwable> result, @Nullable Long 
contenderId, long holderId)
+      throws Exception {
+    String query;
+    switch (backendType) {
+      case "h2":
+        query =
+            "SELECT COUNT(*) FROM INFORMATION_SCHEMA.SESSIONS WHERE BLOCKER_ID 
= "
+                + holderId
+                + (contenderId == null ? "" : " AND SESSION_ID = " + 
contenderId);
+        break;
+      case "mysql":
+        query =
+            "SELECT COUNT(*) FROM performance_schema.data_lock_waits w"
+                + " JOIN performance_schema.threads r ON r.THREAD_ID = 
w.REQUESTING_THREAD_ID"
+                + " JOIN performance_schema.threads b ON b.THREAD_ID = 
w.BLOCKING_THREAD_ID"
+                + " WHERE (b.PROCESSLIST_ID = "
+                + holderId
+                + " OR b.THREAD_ID = r.THREAD_ID)"
+                + (contenderId == null ? "" : " AND r.PROCESSLIST_ID = " + 
contenderId);
+        break;
+      case "postgresql":
+        query =
+            "SELECT COUNT(*) FROM pg_stat_activity a WHERE "
+                + holderId
+                + " = ANY(pg_blocking_pids(a.pid))"
+                + (contenderId == null ? "" : " AND a.pid = " + contenderId);
+        break;
+      default:
+        throw new IllegalStateException("Unsupported backend: " + backendType);
+    }
+    // Observe the actual waiter/blocker pair. A slow thread or connection 
checkout alone cannot
+    // satisfy this assertion, and an unexpectedly completed operation fails 
immediately.
+    try (SqlSession observer =
+        
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true)) 
{
+      Connection connection = observer.getConnection();
+      await()
+          .pollInSameThread()
+          .atMost(10, TimeUnit.SECONDS)
+          .pollInterval(10, TimeUnit.MILLISECONDS)
+          .until(
+              () -> {
+                if (result.isDone()) {
+                  throw new AssertionError(
+                      "Operation completed without waiting for the holder", 
result.get());
+                }
+                try (Statement statement = connection.createStatement();
+                    ResultSet rows = statement.executeQuery(query)) {
+                  return rows.next() && rows.getLong(1) > 0;
+                }
+              });
+    }
+  }
+}
diff --git a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql 
b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
index 947bef162d..2b1e63bf06 100644
--- a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
+++ b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
@@ -101,3 +101,16 @@ CREATE TABLE IF NOT EXISTS `semantic_model_version_info` (
     KEY `idx_smvi_cid` (`catalog_id`),
     KEY `idx_smvi_sid` (`schema_id`)
 ) ENGINE=InnoDB COMMENT 'semantic model version information';
+
+-- Merge duplicate live owners left by concurrent assignments: the newest live 
row
+-- (largest id) wins, and older ones are soft-deleted.
+UPDATE `owner_meta`
+    SET `deleted_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND FROM 
CURRENT_TIMESTAMP(3)) / 1000),
+        `updated_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND FROM 
CURRENT_TIMESTAMP(3)) / 1000)
+    WHERE `deleted_at` = 0
+      AND `id` < (
+        SELECT MAX(d.`id`) FROM `owner_meta` d
+        WHERE d.`deleted_at` = 0
+          AND d.`metadata_object_id` = `owner_meta`.`metadata_object_id`
+          AND d.`metadata_object_type` = `owner_meta`.`metadata_object_type`
+      );
diff --git a/scripts/mysql/schema-2.0.0-mysql.sql 
b/scripts/mysql/schema-2.0.0-mysql.sql
index 70505d3a9e..68e9f9cf23 100644
--- a/scripts/mysql/schema-2.0.0-mysql.sql
+++ b/scripts/mysql/schema-2.0.0-mysql.sql
@@ -325,7 +325,7 @@ CREATE TABLE IF NOT EXISTS `owner_meta` (
     `deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'owner 
relation deleted at',
     `updated_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'updated at',
     PRIMARY KEY (`id`),
-    UNIQUE KEY `uk_ow_me_del` (`owner_id`, `metadata_object_id`, 
`metadata_object_type`,`deleted_at`),
+    UNIQUE KEY `uk_ow_me_del` (`owner_id`, `metadata_object_id`, 
`metadata_object_type`, `deleted_at`),
     KEY `idx_oid` (`owner_id`),
     KEY `idx_meid` (`metadata_object_id`),
     KEY `idx_owner_meta_del_upd_obj` (`deleted_at`, `updated_at`, 
`metadata_object_id`)
diff --git a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql 
b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
index 678767d88d..4bad49e12d 100644
--- a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
+++ b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
@@ -182,3 +182,18 @@ CREATE TABLE IF NOT EXISTS `semantic_model_version_info` (
     KEY `idx_smvi_cid` (`catalog_id`),
     KEY `idx_smvi_sid` (`schema_id`)
 ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin COMMENT 'semantic 
model version information';
+
+-- Merge duplicate live owners left by concurrent assignments: the newest live 
row
+-- (largest id) wins, and older ones are soft-deleted.
+UPDATE `owner_meta` o
+    JOIN (
+        SELECT `metadata_object_id`, `metadata_object_type`, MAX(`id`) AS 
keep_id
+        FROM `owner_meta`
+        WHERE `deleted_at` = 0
+        GROUP BY `metadata_object_id`, `metadata_object_type`
+        HAVING COUNT(*) > 1
+    ) d ON o.`metadata_object_id` = d.`metadata_object_id`
+       AND o.`metadata_object_type` = d.`metadata_object_type`
+    SET o.`deleted_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND 
FROM CURRENT_TIMESTAMP(3)) / 1000),
+        o.`updated_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND 
FROM CURRENT_TIMESTAMP(3)) / 1000)
+    WHERE o.`deleted_at` = 0 AND o.`id` <> d.keep_id;
diff --git a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql 
b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
index 4a64d43798..774cec6b8e 100644
--- a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
+++ b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
@@ -150,3 +150,16 @@ COMMENT ON COLUMN 
semantic_model_version_info.semantic_model_definition IS 'stru
 COMMENT ON COLUMN semantic_model_version_info.properties IS 'semantic model 
properties snapshot (JSON)';
 COMMENT ON COLUMN semantic_model_version_info.audit_info IS 'semantic model 
version audit info';
 COMMENT ON COLUMN semantic_model_version_info.deleted_at IS 'version deleted 
at';
+
+-- Merge duplicate live owners left by concurrent assignments: the newest live 
row
+-- (largest id) wins, and older ones are soft-deleted.
+UPDATE owner_meta
+    SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS 
BIGINT),
+        updated_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS 
BIGINT)
+    WHERE deleted_at = 0
+      AND id < (
+        SELECT MAX(d.id) FROM owner_meta d
+        WHERE d.deleted_at = 0
+          AND d.metadata_object_id = owner_meta.metadata_object_id
+          AND d.metadata_object_type = owner_meta.metadata_object_type
+      );

Reply via email to