924060929 commented on code in PR #66913:
URL: https://github.com/apache/doris/pull/66913#discussion_r3901090448
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/rewrite/RewriteDataFileExecutor.java:
##########
@@ -59,82 +59,91 @@ public RewriteDataFileExecutor(IcebergExternalTable
dorisTable,
/**
* Execute rewrite for multiple groups concurrently
*/
- public RewriteResult executeGroupsConcurrently(List<RewriteDataGroup>
groups, long targetFileSizeBytes)
+ public RewriteResult executeGroupsConcurrently(List<RewriteDataGroup>
groups, long targetFileSizeBytes,
+ WritableTableLease writableTableLease)
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
+ MvccSnapshot targetSnapshot = new IcebergMvccSnapshot(
+ IcebergUtils.getSnapshotForWritableLease(dorisTable,
writableTableLease));
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);
- }
- } catch (JobException e) {
- throw new UserException("Failed to submit rewrite tasks: " +
e.getMessage(), e);
- }
-
- // Wait for all tasks to complete
- waitForTasksCompletion(resultCollector, groups.size());
+ transaction.beginRewrite(dorisTable,
writableTableLease.getTable(), writableTableLease);
- // Finish rewrite operation
- transaction.finishRewrite();
+ // Register files to delete
+ for (RewriteDataGroup group : groups) {
+
transaction.updateRewriteFiles(Lists.newArrayList(group.getDataFiles()));
+ }
- // 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();
+ // Create result collector and 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,
+ writableTableLease.retain(),
+ 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);
+ // Submit tasks to TransientTaskManager
+ 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);
+ // 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();
+
+ transaction.commit();
Review Comment:
Fixed in f0189d2b421. RewriteDataFileExecutor now captures the exact
TransactionManager once, uses it for begin/get/commit/rollback, and places
every post-begin step (including writable snapshot setup) inside the rollback
scope. Successful commit and every failure path therefore remove the
transaction from both the catalog-local and GlobalExternalTransactionInfoMgr
registries. Added success and post-begin setup-failure tests that assert both
registries are empty, plus commit-before-cache-invalidation ordering.
Rewrite/transaction focused FE UT (47 tests), Checkstyle, and ./build.sh --fe
pass.
--
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]