Gabriel39 commented on code in PR #66825:
URL: https://github.com/apache/doris/pull/66825#discussion_r3800770697


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/rewrite/RewriteDataFileExecutor.java:
##########
@@ -61,88 +63,103 @@ public RewriteDataFileExecutor(IcebergExternalTable 
dorisTable,
      */
     public RewriteResult executeGroupsConcurrently(List<RewriteDataGroup> 
groups, long targetFileSizeBytes)
             throws UserException {
-        // Begin transaction
-        long transactionId = 
dorisTable.getCatalog().getTransactionManager().begin();
-        IcebergTransaction transaction = (IcebergTransaction) 
dorisTable.getCatalog().getTransactionManager()
-                .getTransaction(transactionId);
-        MvccSnapshot targetSnapshot = 
dorisTable.loadSnapshot(Optional.empty(), Optional.empty());
-        Table targetIcebergTable = ((IcebergMvccSnapshot) 
targetSnapshot).getSnapshotCacheValue()
-                .getIcebergTable().orElseThrow(
-                        () -> new UserException("Iceberg rewrite target 
metadata is not available"));
-        transaction.beginRewrite(dorisTable, targetIcebergTable);
-
-        // Register files to delete
-        for (RewriteDataGroup group : groups) {
-            
transaction.updateRewriteFiles(Lists.newArrayList(group.getDataFiles()));
-        }
-
-        // Create result collector and tasks
+        TransactionManager transactionManager = 
dorisTable.getCatalog().getTransactionManager();
+        long transactionId = transactionManager.begin();
         List<RewriteGroupTask> tasks = Lists.newArrayList();
-        RewriteResultCollector resultCollector = new 
RewriteResultCollector(groups.size(), tasks);
-
-        // Get available BE count once before creating tasks
-        // This avoids calling getBackendsNumber() in each task during 
multi-threaded execution.
-        // Use compute group from connect context to align with actual BE 
selection for queries.
-        int availableBeCount = getAvailableBeCount();
-
-        // Create tasks with callbacks
-        for (RewriteDataGroup group : groups) {
-            RewriteGroupTask task = new RewriteGroupTask(
-                    group,
-                    transactionId,
-                    dorisTable,
-                    targetSnapshot,
-                    connectContext,
-                    targetFileSizeBytes,
-                    availableBeCount,
-                    new RewriteGroupTask.RewriteResultCallback() {
-                        @Override
-                        public void onTaskCompleted(Long taskId) {
-                            resultCollector.onTaskCompleted(taskId);
-                        }
-
-                        @Override
-                        public void onTaskFailed(Long taskId, Exception error) 
{
-                            resultCollector.onTaskFailed(taskId, error);
-                        }
-                    });
-            tasks.add(task);
-        }
-
-        // Submit tasks to TransientTaskManager
+        boolean committed = false;
         try {
-            for (TransientTaskExecutor task : tasks) {
-                
Env.getCurrentEnv().getTransientTaskManager().addMemoryTask(task);
+            IcebergTransaction transaction = (IcebergTransaction) 
transactionManager
+                    .getTransaction(transactionId);
+            MvccSnapshot targetSnapshot = 
dorisTable.loadSnapshot(Optional.empty(), Optional.empty());
+            Table targetIcebergTable = ((IcebergMvccSnapshot) 
targetSnapshot).getSnapshotCacheValue()
+                    .getIcebergTable().orElseThrow(
+                            () -> new UserException("Iceberg rewrite target 
metadata is not available"));
+            transaction.beginRewrite(dorisTable, targetIcebergTable);
+
+            for (RewriteDataGroup group : groups) {
+                
transaction.updateRewriteFiles(Lists.newArrayList(group.getDataFiles()));
             }
-        } catch (JobException e) {
-            throw new UserException("Failed to submit rewrite tasks: " + 
e.getMessage(), e);
-        }
 
-        // Wait for all tasks to complete
-        waitForTasksCompletion(resultCollector, groups.size());
-
-        // Finish rewrite operation
-        transaction.finishRewrite();
-
-        // Collect statistics from transaction after all tasks are completed
-        int rewrittenDataFilesCount = groups.stream().mapToInt(group -> 
group.getDataFiles().size()).sum();
-        // this should after finishRewrite
-        int addedDataFilesCount = transaction.getFilesToAddCount();
-        long rewrittenBytesCount = groups.stream().mapToLong(group -> 
group.getTotalSize()).sum();
-        int removedDeleteFilesCount = groups.stream().mapToInt(group -> 
group.getDeleteFileCount()).sum();
+            RewriteResultCollector resultCollector = new 
RewriteResultCollector(groups.size(), tasks);
+            int availableBeCount = getAvailableBeCount();
+            for (RewriteDataGroup group : groups) {
+                RewriteGroupTask task = new RewriteGroupTask(
+                        group, transactionId, dorisTable, targetSnapshot, 
connectContext,
+                        targetFileSizeBytes, availableBeCount,
+                        new RewriteGroupTask.RewriteResultCallback() {
+                            @Override
+                            public void onTaskCompleted(Long taskId) {
+                                resultCollector.onTaskCompleted(taskId);
+                            }
+
+                            @Override
+                            public void onTaskFailed(Long taskId, Exception 
error) {
+                                resultCollector.onTaskFailed(taskId, error);
+                            }
+                        });
+                tasks.add(task);
+            }
 
-        commitAndInvalidate(transaction);
+            try {
+                for (TransientTaskExecutor task : tasks) {
+                    
Env.getCurrentEnv().getTransientTaskManager().addMemoryTask(task);
+                }
+            } catch (JobException e) {
+                throw new UserException("Failed to submit rewrite tasks: " + 
e.getMessage(), e);
+            }
 
-        return new RewriteResult(rewrittenDataFilesCount, addedDataFilesCount,
-                rewrittenBytesCount, removedDeleteFilesCount);
+            waitForTasksCompletion(resultCollector, groups.size());
+            transaction.finishRewrite();
+
+            int rewrittenDataFilesCount = groups.stream()
+                    .mapToInt(group -> group.getDataFiles().size()).sum();
+            int addedDataFilesCount = transaction.getFilesToAddCount();
+            long rewrittenBytesCount = groups.stream().mapToLong(group -> 
group.getTotalSize()).sum();
+            int removedDeleteFilesCount = groups.stream()
+                    .mapToInt(group -> group.getDeleteFileCount()).sum();
+
+            commitAndInvalidate(transactionManager, transactionId);

Review Comment:
   Fixed. The transaction is now marked committed immediately after the manager 
commit, before best-effort cache invalidation, so an invalidation failure 
cannot trigger rollback or make the durable rewrite retryable. Added 
failure-injection coverage.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/rewrite/RewriteDataFileExecutor.java:
##########
@@ -61,88 +63,103 @@ public RewriteDataFileExecutor(IcebergExternalTable 
dorisTable,
      */
     public RewriteResult executeGroupsConcurrently(List<RewriteDataGroup> 
groups, long targetFileSizeBytes)
             throws UserException {
-        // Begin transaction
-        long transactionId = 
dorisTable.getCatalog().getTransactionManager().begin();
-        IcebergTransaction transaction = (IcebergTransaction) 
dorisTable.getCatalog().getTransactionManager()
-                .getTransaction(transactionId);
-        MvccSnapshot targetSnapshot = 
dorisTable.loadSnapshot(Optional.empty(), Optional.empty());
-        Table targetIcebergTable = ((IcebergMvccSnapshot) 
targetSnapshot).getSnapshotCacheValue()
-                .getIcebergTable().orElseThrow(
-                        () -> new UserException("Iceberg rewrite target 
metadata is not available"));
-        transaction.beginRewrite(dorisTable, targetIcebergTable);
-
-        // Register files to delete
-        for (RewriteDataGroup group : groups) {
-            
transaction.updateRewriteFiles(Lists.newArrayList(group.getDataFiles()));
-        }
-
-        // Create result collector and tasks
+        TransactionManager transactionManager = 
dorisTable.getCatalog().getTransactionManager();
+        long transactionId = transactionManager.begin();
         List<RewriteGroupTask> tasks = Lists.newArrayList();
-        RewriteResultCollector resultCollector = new 
RewriteResultCollector(groups.size(), tasks);
-
-        // Get available BE count once before creating tasks
-        // This avoids calling getBackendsNumber() in each task during 
multi-threaded execution.
-        // Use compute group from connect context to align with actual BE 
selection for queries.
-        int availableBeCount = getAvailableBeCount();
-
-        // Create tasks with callbacks
-        for (RewriteDataGroup group : groups) {
-            RewriteGroupTask task = new RewriteGroupTask(
-                    group,
-                    transactionId,
-                    dorisTable,
-                    targetSnapshot,
-                    connectContext,
-                    targetFileSizeBytes,
-                    availableBeCount,
-                    new RewriteGroupTask.RewriteResultCallback() {
-                        @Override
-                        public void onTaskCompleted(Long taskId) {
-                            resultCollector.onTaskCompleted(taskId);
-                        }
-
-                        @Override
-                        public void onTaskFailed(Long taskId, Exception error) 
{
-                            resultCollector.onTaskFailed(taskId, error);
-                        }
-                    });
-            tasks.add(task);
-        }
-
-        // Submit tasks to TransientTaskManager
+        boolean committed = false;
         try {
-            for (TransientTaskExecutor task : tasks) {
-                
Env.getCurrentEnv().getTransientTaskManager().addMemoryTask(task);
+            IcebergTransaction transaction = (IcebergTransaction) 
transactionManager
+                    .getTransaction(transactionId);
+            MvccSnapshot targetSnapshot = 
dorisTable.loadSnapshot(Optional.empty(), Optional.empty());
+            Table targetIcebergTable = ((IcebergMvccSnapshot) 
targetSnapshot).getSnapshotCacheValue()
+                    .getIcebergTable().orElseThrow(
+                            () -> new UserException("Iceberg rewrite target 
metadata is not available"));
+            transaction.beginRewrite(dorisTable, targetIcebergTable);
+
+            for (RewriteDataGroup group : groups) {
+                
transaction.updateRewriteFiles(Lists.newArrayList(group.getDataFiles()));
             }
-        } catch (JobException e) {
-            throw new UserException("Failed to submit rewrite tasks: " + 
e.getMessage(), e);
-        }
 
-        // Wait for all tasks to complete
-        waitForTasksCompletion(resultCollector, groups.size());
-
-        // Finish rewrite operation
-        transaction.finishRewrite();
-
-        // Collect statistics from transaction after all tasks are completed
-        int rewrittenDataFilesCount = groups.stream().mapToInt(group -> 
group.getDataFiles().size()).sum();
-        // this should after finishRewrite
-        int addedDataFilesCount = transaction.getFilesToAddCount();
-        long rewrittenBytesCount = groups.stream().mapToLong(group -> 
group.getTotalSize()).sum();
-        int removedDeleteFilesCount = groups.stream().mapToInt(group -> 
group.getDeleteFileCount()).sum();
+            RewriteResultCollector resultCollector = new 
RewriteResultCollector(groups.size(), tasks);
+            int availableBeCount = getAvailableBeCount();
+            for (RewriteDataGroup group : groups) {
+                RewriteGroupTask task = new RewriteGroupTask(
+                        group, transactionId, dorisTable, targetSnapshot, 
connectContext,
+                        targetFileSizeBytes, availableBeCount,
+                        new RewriteGroupTask.RewriteResultCallback() {
+                            @Override
+                            public void onTaskCompleted(Long taskId) {
+                                resultCollector.onTaskCompleted(taskId);
+                            }
+
+                            @Override
+                            public void onTaskFailed(Long taskId, Exception 
error) {
+                                resultCollector.onTaskFailed(taskId, error);
+                            }
+                        });
+                tasks.add(task);
+            }
 
-        commitAndInvalidate(transaction);
+            try {
+                for (TransientTaskExecutor task : tasks) {
+                    
Env.getCurrentEnv().getTransientTaskManager().addMemoryTask(task);
+                }
+            } catch (JobException e) {
+                throw new UserException("Failed to submit rewrite tasks: " + 
e.getMessage(), e);
+            }
 
-        return new RewriteResult(rewrittenDataFilesCount, addedDataFilesCount,
-                rewrittenBytesCount, removedDeleteFilesCount);
+            waitForTasksCompletion(resultCollector, groups.size());

Review Comment:
   Fixed. Iceberg rewrite group executors no longer commit or roll back the 
parent transaction, and group failures propagate to the parent collector before 
source deletion. Added tests for no per-group commit and failure propagation.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to