This is an automated email from the ASF dual-hosted git repository.
danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new a8e9db446c3 [HUDI-7530] Refactoring of handleUpdateInternal in
CommitActionExecutors and HoodieTables (#10908)
a8e9db446c3 is described below
commit a8e9db446c362a364a49749fab795be31fc33afb
Author: wombatu-kun <[email protected]>
AuthorDate: Sat Mar 23 08:07:29 2024 +0700
[HUDI-7530] Refactoring of handleUpdateInternal in CommitActionExecutors
and HoodieTables (#10908)
Co-authored-by: Vova Kolmakov <[email protected]>
---
.../org/apache/hudi/io/HoodieAppendHandle.java | 2 +-
.../java/org/apache/hudi/io/HoodieMergeHandle.java | 9 ++++++++
.../java/org/apache/hudi/io/HoodieWriteHandle.java | 2 +-
.../java/org/apache/hudi/table/HoodieTable.java | 10 ++++++++
.../hudi/table/HoodieFlinkCopyOnWriteTable.java | 18 ++-------------
.../commit/BaseFlinkCommitActionExecutor.java | 16 ++-----------
.../hudi/table/HoodieJavaCopyOnWriteTable.java | 18 ++-------------
.../commit/BaseJavaCommitActionExecutor.java | 14 ++---------
.../hudi/table/HoodieSparkCopyOnWriteTable.java | 27 ++--------------------
.../org/apache/hudi/table/HoodieSparkTable.java | 22 ++++++++++++++++++
.../bootstrap/BaseBootstrapMetadataHandler.java | 2 +-
.../commit/BaseSparkCommitActionExecutor.java | 26 ++-------------------
12 files changed, 56 insertions(+), 110 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
index 93df86e170d..1301f046ae2 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
@@ -562,7 +562,7 @@ public class HoodieAppendHandle<T, I, K, O> extends
HoodieWriteHandle<T, I, K, O
return IOType.APPEND;
}
- public List<WriteStatus> writeStatuses() {
+ public List<WriteStatus> getWriteStatuses() {
return statuses;
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandle.java
index d882a68e17c..4f5f240c4fd 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandle.java
@@ -482,6 +482,15 @@ public class HoodieMergeHandle<T, I, K, O> extends
HoodieWriteHandle<T, I, K, O>
}
}
+ public Iterator<List<WriteStatus>> getWriteStatusesAsIterator() {
+ List<WriteStatus> statuses = getWriteStatuses();
+ // TODO(vc): This needs to be revisited
+ if (getPartitionPath() == null) {
+ LOG.info("Upsert Handle has partition path as null {}, {}",
getOldFilePath(), statuses);
+ }
+ return Collections.singletonList(statuses).iterator();
+ }
+
public Path getOldFilePath() {
return oldFilePath;
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteHandle.java
index 9d1bb6d511e..ab80629c941 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteHandle.java
@@ -190,7 +190,7 @@ public abstract class HoodieWriteHandle<T, I, K, O> extends
HoodieIOHandle<T, I,
public abstract List<WriteStatus> close();
- public List<WriteStatus> writeStatuses() {
+ public List<WriteStatus> getWriteStatuses() {
return Collections.singletonList(writeStatus);
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
index 36990b1e9e8..aadb0d48685 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
@@ -71,11 +71,13 @@ import
org.apache.hudi.exception.SchemaCompatibilityException;
import org.apache.hudi.hadoop.fs.ConsistencyGuard;
import org.apache.hudi.hadoop.fs.ConsistencyGuard.FileVisibility;
import org.apache.hudi.index.HoodieIndex;
+import org.apache.hudi.io.HoodieMergeHandle;
import org.apache.hudi.metadata.HoodieTableMetadata;
import org.apache.hudi.metadata.HoodieTableMetadataWriter;
import org.apache.hudi.metadata.MetadataPartitionType;
import org.apache.hudi.table.action.HoodieWriteMetadata;
import org.apache.hudi.table.action.bootstrap.HoodieBootstrapWriteMetadata;
+import org.apache.hudi.table.action.commit.HoodieMergeHelper;
import org.apache.hudi.table.marker.WriteMarkers;
import org.apache.hudi.table.marker.WriteMarkersFactory;
import org.apache.hudi.table.storage.HoodieLayoutFactory;
@@ -1093,4 +1095,12 @@ public abstract class HoodieTable<T, I, K, O> implements
Serializable {
}
return new HashSet<>(Arrays.asList(partitionFields.get()));
}
+
+ public void runMerge(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String
instantTime, String fileId) throws IOException {
+ if (upsertHandle.getOldFilePath() == null) {
+ throw new HoodieUpsertException("Error in finding the old file path at
commit " + instantTime + " for fileId: " + fileId);
+ } else {
+ HoodieMergeHelper.newInstance().runMerge(this, upsertHandle);
+ }
+ }
}
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkCopyOnWriteTable.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkCopyOnWriteTable.java
index c870d736a3d..05779339780 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkCopyOnWriteTable.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkCopyOnWriteTable.java
@@ -41,7 +41,6 @@ import org.apache.hudi.common.util.Option;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.exception.HoodieIOException;
import org.apache.hudi.exception.HoodieNotSupportedException;
-import org.apache.hudi.exception.HoodieUpsertException;
import org.apache.hudi.io.HoodieCreateHandle;
import org.apache.hudi.io.HoodieMergeHandle;
import org.apache.hudi.io.HoodieMergeHandleFactory;
@@ -64,7 +63,6 @@ import
org.apache.hudi.table.action.commit.FlinkInsertOverwriteTableCommitAction
import
org.apache.hudi.table.action.commit.FlinkInsertPreppedCommitActionExecutor;
import org.apache.hudi.table.action.commit.FlinkUpsertCommitActionExecutor;
import
org.apache.hudi.table.action.commit.FlinkUpsertPreppedCommitActionExecutor;
-import org.apache.hudi.table.action.commit.HoodieMergeHelper;
import org.apache.hudi.table.action.rollback.BaseRollbackPlanActionExecutor;
import org.apache.hudi.table.action.rollback.CopyOnWriteRollbackActionExecutor;
import org.slf4j.Logger;
@@ -421,20 +419,8 @@ public class HoodieFlinkCopyOnWriteTable<T>
protected Iterator<List<WriteStatus>>
handleUpdateInternal(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String
instantTime,
String fileId)
throws IOException {
- if (upsertHandle.getOldFilePath() == null) {
- throw new HoodieUpsertException(
- "Error in finding the old file path at commit " + instantTime + "
for fileId: " + fileId);
- } else {
- HoodieMergeHelper.newInstance().runMerge(this, upsertHandle);
- }
-
- // TODO(vc): This needs to be revisited
- if (upsertHandle.getPartitionPath() == null) {
- LOG.info("Upsert Handle has partition path as null " +
upsertHandle.getOldFilePath() + ", "
- + upsertHandle.writeStatuses());
- }
-
- return Collections.singletonList(upsertHandle.writeStatuses()).iterator();
+ runMerge(upsertHandle, instantTime, fileId);
+ return upsertHandle.getWriteStatusesAsIterator();
}
protected HoodieMergeHandle getUpdateHandle(String instantTime, String
partitionPath, String fileId,
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/BaseFlinkCommitActionExecutor.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/BaseFlinkCommitActionExecutor.java
index 732832ae911..750ef71bc37 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/BaseFlinkCommitActionExecutor.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/BaseFlinkCommitActionExecutor.java
@@ -184,20 +184,8 @@ public abstract class BaseFlinkCommitActionExecutor<T>
extends
protected Iterator<List<WriteStatus>>
handleUpdateInternal(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String fileId)
throws IOException {
- if (upsertHandle.getOldFilePath() == null) {
- throw new HoodieUpsertException(
- "Error in finding the old file path at commit " + instantTime + "
for fileId: " + fileId);
- } else {
- HoodieMergeHelper.newInstance().runMerge(table, upsertHandle);
- }
-
- // TODO(vc): This needs to be revisited
- if (upsertHandle.getPartitionPath() == null) {
- LOG.info("Upsert Handle has partition path as null " +
upsertHandle.getOldFilePath() + ", "
- + upsertHandle.writeStatuses());
- }
-
- return Collections.singletonList(upsertHandle.writeStatuses()).iterator();
+ table.runMerge(upsertHandle, instantTime, fileId);
+ return upsertHandle.getWriteStatusesAsIterator();
}
@Override
diff --git
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/HoodieJavaCopyOnWriteTable.java
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/HoodieJavaCopyOnWriteTable.java
index 6215a3aea8f..6e111f67da2 100644
---
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/HoodieJavaCopyOnWriteTable.java
+++
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/HoodieJavaCopyOnWriteTable.java
@@ -42,7 +42,6 @@ import org.apache.hudi.common.util.Option;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.exception.HoodieIOException;
import org.apache.hudi.exception.HoodieNotSupportedException;
-import org.apache.hudi.exception.HoodieUpsertException;
import org.apache.hudi.io.HoodieCreateHandle;
import org.apache.hudi.io.HoodieMergeHandle;
import org.apache.hudi.io.HoodieMergeHandleFactory;
@@ -55,7 +54,6 @@ import org.apache.hudi.table.action.clean.CleanActionExecutor;
import org.apache.hudi.table.action.clean.CleanPlanActionExecutor;
import org.apache.hudi.table.action.cluster.ClusteringPlanActionExecutor;
import
org.apache.hudi.table.action.cluster.JavaExecuteClusteringCommitActionExecutor;
-import org.apache.hudi.table.action.commit.HoodieMergeHelper;
import org.apache.hudi.table.action.commit.JavaBulkInsertCommitActionExecutor;
import
org.apache.hudi.table.action.commit.JavaBulkInsertPreppedCommitActionExecutor;
import org.apache.hudi.table.action.commit.JavaDeleteCommitActionExecutor;
@@ -290,20 +288,8 @@ public class HoodieJavaCopyOnWriteTable<T>
protected Iterator<List<WriteStatus>>
handleUpdateInternal(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String
instantTime,
String fileId)
throws IOException {
- if (upsertHandle.getOldFilePath() == null) {
- throw new HoodieUpsertException(
- "Error in finding the old file path at commit " + instantTime + "
for fileId: " + fileId);
- } else {
- HoodieMergeHelper.newInstance().runMerge(this, upsertHandle);
- }
-
- // TODO(yihua): This needs to be revisited
- if (upsertHandle.getPartitionPath() == null) {
- LOG.info("Upsert Handle has partition path as null " +
upsertHandle.getOldFilePath() + ", "
- + upsertHandle.writeStatuses());
- }
-
- return Collections.singletonList(upsertHandle.writeStatuses()).iterator();
+ runMerge(upsertHandle, instantTime, fileId);
+ return upsertHandle.getWriteStatusesAsIterator();
}
protected HoodieMergeHandle getUpdateHandle(String instantTime, String
partitionPath, String fileId,
diff --git
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java
index 5ce6bcccef5..84e4316d164 100644
---
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java
+++
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java
@@ -243,18 +243,8 @@ public abstract class BaseJavaCommitActionExecutor<T>
extends
protected Iterator<List<WriteStatus>>
handleUpdateInternal(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String fileId)
throws IOException {
- if (upsertHandle.getOldFilePath() == null) {
- throw new HoodieUpsertException(
- "Error in finding the old file path at commit " + instantTime + "
for fileId: " + fileId);
- } else {
- HoodieMergeHelper.newInstance().runMerge(table, upsertHandle);
- }
-
- List<WriteStatus> statuses = upsertHandle.writeStatuses();
- if (upsertHandle.getPartitionPath() == null) {
- LOG.info("Upsert Handle has partition path as null " +
upsertHandle.getOldFilePath() + ", " + statuses);
- }
- return Collections.singletonList(statuses).iterator();
+ table.runMerge(upsertHandle, instantTime, fileId);
+ return upsertHandle.getWriteStatusesAsIterator();
}
protected HoodieMergeHandle getUpdateHandle(String partitionPath, String
fileId, Iterator<HoodieRecord<T>> recordItr) {
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkCopyOnWriteTable.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkCopyOnWriteTable.java
index b9e1033acdd..0a533e65912 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkCopyOnWriteTable.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkCopyOnWriteTable.java
@@ -30,7 +30,6 @@ import org.apache.hudi.avro.model.HoodieRollbackMetadata;
import org.apache.hudi.avro.model.HoodieRollbackPlan;
import org.apache.hudi.avro.model.HoodieSavepointMetadata;
import org.apache.hudi.client.WriteStatus;
-import org.apache.hudi.client.utils.SparkPartitionUtils;
import org.apache.hudi.client.common.HoodieSparkEngineContext;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.data.HoodieData;
@@ -47,7 +46,6 @@ import org.apache.hudi.exception.HoodieException;
import org.apache.hudi.exception.HoodieIOException;
import org.apache.hudi.exception.HoodieMetadataException;
import org.apache.hudi.exception.HoodieNotSupportedException;
-import org.apache.hudi.exception.HoodieUpsertException;
import org.apache.hudi.io.HoodieCreateHandle;
import org.apache.hudi.io.HoodieMergeHandle;
import org.apache.hudi.io.HoodieMergeHandleFactory;
@@ -61,7 +59,6 @@ import org.apache.hudi.table.action.clean.CleanActionExecutor;
import org.apache.hudi.table.action.clean.CleanPlanActionExecutor;
import org.apache.hudi.table.action.cluster.ClusteringPlanActionExecutor;
import
org.apache.hudi.table.action.cluster.SparkExecuteClusteringCommitActionExecutor;
-import org.apache.hudi.table.action.commit.HoodieMergeHelper;
import org.apache.hudi.table.action.commit.SparkBulkInsertCommitActionExecutor;
import
org.apache.hudi.table.action.commit.SparkBulkInsertPreppedCommitActionExecutor;
import org.apache.hudi.table.action.commit.SparkDeleteCommitActionExecutor;
@@ -248,28 +245,8 @@ public class HoodieSparkCopyOnWriteTable<T>
protected Iterator<List<WriteStatus>>
handleUpdateInternal(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String
instantTime,
String fileId)
throws IOException {
- if (upsertHandle.getOldFilePath() == null) {
- throw new HoodieUpsertException(
- "Error in finding the old file path at commit " + instantTime + "
for fileId: " + fileId);
- } else {
- if (upsertHandle.baseFileForMerge().getBootstrapBaseFile().isPresent()) {
- Option<String[]> partitionFields =
getMetaClient().getTableConfig().getPartitionFields();
- Object[] partitionValues =
SparkPartitionUtils.getPartitionFieldVals(partitionFields,
upsertHandle.getPartitionPath(),
- getMetaClient().getTableConfig().getBootstrapBasePath().get(),
- upsertHandle.getWriterSchema(), getHadoopConf());
- upsertHandle.setPartitionFields(partitionFields);
- upsertHandle.setPartitionValues(partitionValues);
- }
- HoodieMergeHelper.newInstance().runMerge(this, upsertHandle);
- }
-
- // TODO(vc): This needs to be revisited
- if (upsertHandle.getPartitionPath() == null) {
- LOG.info("Upsert Handle has partition path as null " +
upsertHandle.getOldFilePath() + ", "
- + upsertHandle.writeStatuses());
- }
-
- return Collections.singletonList(upsertHandle.writeStatuses()).iterator();
+ runMerge(upsertHandle, instantTime, fileId);
+ return upsertHandle.getWriteStatusesAsIterator();
}
protected HoodieMergeHandle getUpdateHandle(String instantTime, String
partitionPath, String fileId,
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
index 08d8a88ae1b..fc3aba63740 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
@@ -19,6 +19,7 @@
package org.apache.hudi.table;
import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.client.utils.SparkPartitionUtils;
import org.apache.hudi.common.data.HoodieData;
import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
@@ -30,12 +31,15 @@ import org.apache.hudi.common.util.Option;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.exception.HoodieException;
import org.apache.hudi.exception.HoodieMetadataException;
+import org.apache.hudi.exception.HoodieUpsertException;
import org.apache.hudi.index.HoodieIndex;
import org.apache.hudi.index.SparkHoodieIndexFactory;
+import org.apache.hudi.io.HoodieMergeHandle;
import org.apache.hudi.metadata.HoodieTableMetadata;
import org.apache.hudi.metadata.HoodieTableMetadataWriter;
import org.apache.hudi.metadata.SparkHoodieBackedTableMetadataWriter;
import org.apache.hadoop.fs.Path;
+import org.apache.hudi.table.action.commit.HoodieMergeHelper;
import org.apache.spark.TaskContext;
import org.apache.spark.TaskContext$;
@@ -125,4 +129,22 @@ public abstract class HoodieSparkTable<T>
final TaskContext taskContext = TaskContext.get();
return () -> TaskContext$.MODULE$.setTaskContext(taskContext);
}
+
+ @Override
+ public void runMerge(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String
instantTime, String fileId) throws IOException {
+ if (upsertHandle.getOldFilePath() == null) {
+ throw new HoodieUpsertException("Error in finding the old file path at
commit " + instantTime + " for fileId: " + fileId);
+ } else {
+ if (upsertHandle.baseFileForMerge().getBootstrapBaseFile().isPresent()) {
+ Option<String[]> partitionFields =
getMetaClient().getTableConfig().getPartitionFields();
+ Object[] partitionValues =
SparkPartitionUtils.getPartitionFieldVals(partitionFields,
upsertHandle.getPartitionPath(),
+ getMetaClient().getTableConfig().getBootstrapBasePath().get(),
+ upsertHandle.getWriterSchema(), getHadoopConf());
+ upsertHandle.setPartitionFields(partitionFields);
+ upsertHandle.setPartitionValues(partitionValues);
+ }
+ HoodieMergeHelper.newInstance().runMerge(this, upsertHandle);
+ }
+ }
+
}
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/bootstrap/BaseBootstrapMetadataHandler.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/bootstrap/BaseBootstrapMetadataHandler.java
index 4d6d07c9e49..ffda89d5b7f 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/bootstrap/BaseBootstrapMetadataHandler.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/bootstrap/BaseBootstrapMetadataHandler.java
@@ -70,7 +70,7 @@ public abstract class BaseBootstrapMetadataHandler implements
BootstrapMetadataH
throw new HoodieException(e.getMessage(), e);
}
- BootstrapWriteStatus writeStatus = (BootstrapWriteStatus)
bootstrapHandle.writeStatuses().get(0);
+ BootstrapWriteStatus writeStatus = (BootstrapWriteStatus)
bootstrapHandle.getWriteStatuses().get(0);
BootstrapFileMapping bootstrapFileMapping = new BootstrapFileMapping(
config.getBootstrapSourceBasePath(), srcPartitionPath, partitionPath,
srcFileStatus, writeStatus.getFileId());
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java
index 5e379d3e956..0fcb2359cdf 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java
@@ -20,7 +20,6 @@ package org.apache.hudi.table.action.commit;
import org.apache.hudi.client.WriteStatus;
import
org.apache.hudi.client.clustering.update.strategy.SparkAllowUpdateStrategy;
-import org.apache.hudi.client.utils.SparkPartitionUtils;
import org.apache.hudi.client.utils.SparkValidatorUtils;
import org.apache.hudi.common.data.HoodieData;
import org.apache.hudi.common.data.HoodieData.HoodieDataCacheKey;
@@ -349,29 +348,8 @@ public abstract class BaseSparkCommitActionExecutor<T>
extends
protected Iterator<List<WriteStatus>>
handleUpdateInternal(HoodieMergeHandle<?, ?, ?, ?> upsertHandle, String fileId)
throws IOException {
- if (upsertHandle.getOldFilePath() == null) {
- throw new HoodieUpsertException(
- "Error in finding the old file path at commit " + instantTime + "
for fileId: " + fileId);
- } else {
- if (upsertHandle.baseFileForMerge().getBootstrapBaseFile().isPresent()) {
- Option<String[]> partitionFields =
table.getMetaClient().getTableConfig().getPartitionFields();
- Object[] partitionValues =
SparkPartitionUtils.getPartitionFieldVals(partitionFields,
upsertHandle.getPartitionPath(),
-
table.getMetaClient().getTableConfig().getBootstrapBasePath().get(),
- upsertHandle.getWriterSchema(), table.getHadoopConf());
- upsertHandle.setPartitionFields(partitionFields);
- upsertHandle.setPartitionValues(partitionValues);
- }
-
- HoodieMergeHelper.newInstance().runMerge(table, upsertHandle);
- }
-
- // TODO(vc): This needs to be revisited
- if (upsertHandle.getPartitionPath() == null) {
- LOG.info("Upsert Handle has partition path as null " +
upsertHandle.getOldFilePath() + ", "
- + upsertHandle.writeStatuses());
- }
-
- return Collections.singletonList(upsertHandle.writeStatuses()).iterator();
+ table.runMerge(upsertHandle, instantTime, fileId);
+ return upsertHandle.getWriteStatusesAsIterator();
}
protected HoodieMergeHandle getUpdateHandle(String partitionPath, String
fileId, Iterator<HoodieRecord<T>> recordItr) {