From 0d7eab54ca8e636e4b013578074873c4a9e6d532 Mon Sep 17 00:00:00 2001 From: sychen Date: Thu, 10 Sep 2026 11:09:56 +0800 Subject: [PATCH] [spark] Close the commit after truncating a table or partitions --- .../paimon/spark/PaimonSparkTableBase.scala | 6 ++- .../TruncatePaimonTableWithFilterExec.scala | 53 ++++++++++--------- 2 files changed, 33 insertions(+), 26 deletions(-) diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala index d42d129925dd..d139ed025ba8 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonSparkTableBase.scala @@ -186,7 +186,11 @@ abstract class PaimonSparkTableBase(val table: Table) def truncateTable: Boolean = { val commit = table.newBatchWriteBuilder().newCommit() - commit.truncateTable() + try { + commit.truncateTable() + } finally { + commit.close() + } true } } diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/TruncatePaimonTableWithFilterExec.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/TruncatePaimonTableWithFilterExec.scala index 41c2187ce37a..3093d15e843e 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/TruncatePaimonTableWithFilterExec.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/TruncatePaimonTableWithFilterExec.scala @@ -37,31 +37,34 @@ case class TruncatePaimonTableWithFilterExec( override def run(): Seq[InternalRow] = { val commit = table.newBatchWriteBuilder().newCommit() - - partitionPredicate match { - case Some(p) => - table match { - case fileStoreTable: FileStoreTable => - val matchedPartitions = - fileStoreTable.newSnapshotReader().withPartitionFilter(p).partitions().asScala - if (matchedPartitions.nonEmpty) { - val partitionComputer = new InternalRowPartitionComputer( - fileStoreTable.coreOptions().partitionDefaultName(), - fileStoreTable.schema().logicalPartitionType(), - fileStoreTable.partitionKeys.asScala.toArray, - fileStoreTable.coreOptions().legacyPartitionName() - ) - val dropPartitions = - matchedPartitions.map(partitionComputer.generatePartValues(_).asScala.asJava) - commit.truncatePartitions(dropPartitions.asJava) - } else { - commit.commit(JCollections.emptyList()) - } - case _ => - throw new UnsupportedOperationException("Unsupported truncate table") - } - case _ => - commit.truncateTable() + try { + partitionPredicate match { + case Some(p) => + table match { + case fileStoreTable: FileStoreTable => + val matchedPartitions = + fileStoreTable.newSnapshotReader().withPartitionFilter(p).partitions().asScala + if (matchedPartitions.nonEmpty) { + val partitionComputer = new InternalRowPartitionComputer( + fileStoreTable.coreOptions().partitionDefaultName(), + fileStoreTable.schema().logicalPartitionType(), + fileStoreTable.partitionKeys.asScala.toArray, + fileStoreTable.coreOptions().legacyPartitionName() + ) + val dropPartitions = + matchedPartitions.map(partitionComputer.generatePartValues(_).asScala.asJava) + commit.truncatePartitions(dropPartitions.asJava) + } else { + commit.commit(JCollections.emptyList()) + } + case _ => + throw new UnsupportedOperationException("Unsupported truncate table") + } + case _ => + commit.truncateTable() + } + } finally { + commit.close() } Nil }