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 9b77eb1f5c6 [HUDI-7879] Optimize the redundant creation of HoodieTable
in DataSourceInternalWriterHelper and the unnecessary parameters in createTable
within BaseHoodieWriteClien (#11456)
9b77eb1f5c6 is described below
commit 9b77eb1f5c6e70d0f22b1ecf70532d1876a59f22
Author: majian <[email protected]>
AuthorDate: Sat Jun 15 08:53:28 2024 +0800
[HUDI-7879] Optimize the redundant creation of HoodieTable in
DataSourceInternalWriterHelper and the unnecessary parameters in createTable
within BaseHoodieWriteClien (#11456)
---
.../apache/hudi/client/BaseHoodieWriteClient.java | 43 +++++++++++-----------
.../apache/hudi/client/HoodieFlinkWriteClient.java | 5 +--
.../apache/hudi/client/HoodieJavaWriteClient.java | 5 +--
.../apache/hudi/client/SparkRDDWriteClient.java | 5 +--
.../internal/DataSourceInternalWriterHelper.java | 4 +-
5 files changed, 28 insertions(+), 34 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
index b9d09c735ea..e1d0524ce6f 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java
@@ -89,7 +89,6 @@ import org.apache.hudi.table.upgrade.UpgradeDowngrade;
import com.codahale.metrics.Timer;
import org.apache.avro.Schema;
-import org.apache.hadoop.conf.Configuration;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -223,7 +222,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
}
LOG.info("Committing " + instantTime + " action " + commitActionType);
// Create a Hoodie table which encapsulated the commits and files visible
- HoodieTable table = createTable(config, hadoopConf);
+ HoodieTable table = createTable(config);
HoodieCommitMetadata metadata = CommitUtils.buildMetadata(stats,
partitionToReplaceFileIds,
extraMetadata, operationType, config.getWriteSchema(),
commitActionType);
HoodieInstant inflightInstant = new HoodieInstant(State.INFLIGHT,
commitActionType, instantTime);
@@ -321,9 +320,9 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
}
}
- protected abstract HoodieTable<T, I, K, O> createTable(HoodieWriteConfig
config, Configuration hadoopConf);
+ protected abstract HoodieTable<T, I, K, O> createTable(HoodieWriteConfig
config);
- protected abstract HoodieTable<T, I, K, O> createTable(HoodieWriteConfig
config, Configuration hadoopConf, HoodieTableMetaClient metaClient);
+ protected abstract HoodieTable<T, I, K, O> createTable(HoodieWriteConfig
config, HoodieTableMetaClient metaClient);
void emitCommitMetrics(String instantTime, HoodieCommitMetadata metadata,
String actionType) {
if (writeTimer != null) {
@@ -345,7 +344,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
protected void preCommit(HoodieInstant inflightInstant, HoodieCommitMetadata
metadata) {
// Create a Hoodie table after startTxn which encapsulated the commits and
files visible.
// Important to create this after the lock to ensure the latest commits
show up in the timeline without need for reload
- HoodieTable table = createTable(config, hadoopConf);
+ HoodieTable table = createTable(config);
resolveWriteConflict(table, metadata,
this.pendingInflightAndRequestedInstants);
}
@@ -594,14 +593,14 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* Run any pending compactions.
*/
public void runAnyPendingCompactions() {
- tableServiceClient.runAnyPendingCompactions(createTable(config,
hadoopConf));
+ tableServiceClient.runAnyPendingCompactions(createTable(config));
}
/**
* Run any pending log compactions.
*/
public void runAnyPendingLogCompactions() {
- tableServiceClient.runAnyPendingLogCompactions(createTable(config,
hadoopConf));
+ tableServiceClient.runAnyPendingLogCompactions(createTable(config));
}
/**
@@ -611,7 +610,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @param comment Comment for the savepoint
*/
public void savepoint(String user, String comment) {
- HoodieTable<T, I, K, O> table = createTable(config, hadoopConf);
+ HoodieTable<T, I, K, O> table = createTable(config);
if (table.getCompletedCommitsTimeline().empty()) {
throw new HoodieSavepointException("Could not savepoint. Commit timeline
is empty");
}
@@ -635,7 +634,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @param comment Comment for the savepoint
*/
public void savepoint(String instantTime, String user, String comment) {
- HoodieTable<T, I, K, O> table = createTable(config, hadoopConf);
+ HoodieTable<T, I, K, O> table = createTable(config);
table.savepoint(context, instantTime, user, comment);
}
@@ -643,7 +642,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* Delete a savepoint based on the latest commit action on the savepoint
timeline.
*/
public void deleteSavepoint() {
- HoodieTable<T, I, K, O> table = createTable(config, hadoopConf);
+ HoodieTable<T, I, K, O> table = createTable(config);
HoodieTimeline savePointTimeline =
table.getActiveTimeline().getSavePointTimeline();
if (savePointTimeline.empty()) {
throw new HoodieSavepointException("Could not delete savepoint.
Savepoint timeline is empty");
@@ -661,7 +660,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @param savepointTime Savepoint time to delete
*/
public void deleteSavepoint(String savepointTime) {
- HoodieTable<T, I, K, O> table = createTable(config, hadoopConf);
+ HoodieTable<T, I, K, O> table = createTable(config);
SavepointHelpers.deleteSavepoint(table, savepointTime);
}
@@ -669,7 +668,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* Restore the data to a savepoint based on the latest commit action on the
savepoint timeline.
*/
public void restoreToSavepoint() {
- HoodieTable<T, I, K, O> table = createTable(config, hadoopConf);
+ HoodieTable<T, I, K, O> table = createTable(config);
HoodieTimeline savePointTimeline =
table.getActiveTimeline().getSavePointTimeline();
if (savePointTimeline.empty()) {
throw new HoodieSavepointException("Could not restore to savepoint.
Savepoint timeline is empty");
@@ -866,7 +865,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
*/
public void archive() {
// Create a Hoodie table which encapsulated the commits and files visible
- HoodieTable table = createTable(config, hadoopConf);
+ HoodieTable table = createTable(config);
archive(table);
}
@@ -963,7 +962,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
*/
public Option<String> scheduleIndexing(List<MetadataPartitionType>
partitionTypes, List<String> partitionPaths) {
String instantTime = createNewInstantTime();
- Option<HoodieIndexPlan> indexPlan = createTable(config, hadoopConf)
+ Option<HoodieIndexPlan> indexPlan = createTable(config)
.scheduleIndexing(context, instantTime, partitionTypes,
partitionPaths);
return indexPlan.isPresent() ? Option.of(instantTime) : Option.empty();
}
@@ -975,7 +974,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @return {@link Option<HoodieIndexCommitMetadata>} after successful
indexing.
*/
public Option<HoodieIndexCommitMetadata> index(String indexInstantTime) {
- return createTable(config, hadoopConf).index(context, indexInstantTime);
+ return createTable(config).index(context, indexInstantTime);
}
/**
@@ -984,7 +983,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @param partitionTypes - list of {@link MetadataPartitionType} which needs
to be indexed
*/
public void dropIndex(List<MetadataPartitionType> partitionTypes) {
- HoodieTable table = createTable(config, hadoopConf);
+ HoodieTable table = createTable(config);
String dropInstant = createNewInstantTime();
HoodieInstant ownerInstant = new HoodieInstant(true,
HoodieTimeline.INDEXING_ACTION, dropInstant);
this.txnManager.beginTransaction(Option.of(ownerInstant), Option.empty());
@@ -1089,7 +1088,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
*/
public void commitLogCompaction(String logCompactionInstantTime,
HoodieCommitMetadata metadata,
Option<Map<String, String>> extraMetadata) {
- HoodieTable table = createTable(config,
context.getStorageConf().unwrapAs(Configuration.class));
+ HoodieTable table = createTable(config);
extraMetadata.ifPresent(m -> m.forEach(metadata::addMetadata));
completeLogCompaction(metadata, table, logCompactionInstantTime);
}
@@ -1108,7 +1107,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @return Collection of Write Status
*/
protected HoodieWriteMetadata<O> compact(String compactionInstantTime,
boolean shouldComplete) {
- HoodieTable table = createTable(config,
context.getStorageConf().unwrapAs(Configuration.class));
+ HoodieTable table = createTable(config);
preWrite(compactionInstantTime, WriteOperationType.COMPACT,
table.getMetaClient());
return tableServiceClient.compact(compactionInstantTime, shouldComplete);
}
@@ -1129,7 +1128,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @return Collection of Write Status
*/
protected HoodieWriteMetadata<O> logCompact(String logCompactionInstantTime,
boolean shouldComplete) {
- HoodieTable table = createTable(config,
context.getStorageConf().unwrapAs(Configuration.class));
+ HoodieTable table = createTable(config);
preWrite(logCompactionInstantTime, WriteOperationType.LOG_COMPACT,
table.getMetaClient());
return tableServiceClient.logCompact(logCompactionInstantTime,
shouldComplete);
}
@@ -1167,13 +1166,13 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
* @return Collection of Write Status
*/
public HoodieWriteMetadata<O> cluster(String clusteringInstant, boolean
shouldComplete) {
- HoodieTable table = createTable(config,
context.getStorageConf().unwrapAs(Configuration.class));
+ HoodieTable table = createTable(config);
preWrite(clusteringInstant, WriteOperationType.CLUSTER,
table.getMetaClient());
return tableServiceClient.cluster(clusteringInstant, shouldComplete);
}
public boolean purgePendingClustering(String clusteringInstant) {
- HoodieTable table = createTable(config,
context.getStorageConf().unwrapAs(Configuration.class));
+ HoodieTable table = createTable(config);
preWrite(clusteringInstant, WriteOperationType.CLUSTER,
table.getMetaClient());
return tableServiceClient.purgePendingClustering(clusteringInstant);
}
@@ -1269,7 +1268,7 @@ public abstract class BaseHoodieWriteClient<T, I, K, O>
extends BaseHoodieClient
}
doInitTable(operationType, metaClient, instantTime);
- HoodieTable table = createTable(config, hadoopConf, metaClient);
+ HoodieTable table = createTable(config, metaClient);
// Validate table properties
metaClient.validateTableProperties(config.getProps());
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
index 1d6bf73b2ec..d6b1730566f 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java
@@ -51,7 +51,6 @@ import org.apache.hudi.table.upgrade.UpgradeDowngrade;
import org.apache.hudi.util.WriteStatMerger;
import com.codahale.metrics.Timer;
-import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -114,12 +113,12 @@ public class HoodieFlinkWriteClient<T> extends
}
@Override
- protected HoodieTable createTable(HoodieWriteConfig config, Configuration
hadoopConf) {
+ protected HoodieTable createTable(HoodieWriteConfig config) {
return HoodieFlinkTable.create(config, (HoodieFlinkEngineContext) context);
}
@Override
- protected HoodieTable createTable(HoodieWriteConfig config, Configuration
hadoopConf, HoodieTableMetaClient metaClient) {
+ protected HoodieTable createTable(HoodieWriteConfig config,
HoodieTableMetaClient metaClient) {
return HoodieFlinkTable.create(config, (HoodieFlinkEngineContext) context,
metaClient);
}
diff --git
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/HoodieJavaWriteClient.java
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/HoodieJavaWriteClient.java
index 596767e8cc6..9552804ac97 100644
---
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/HoodieJavaWriteClient.java
+++
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/HoodieJavaWriteClient.java
@@ -43,7 +43,6 @@ import org.apache.hudi.table.action.HoodieWriteMetadata;
import org.apache.hudi.table.upgrade.JavaUpgradeDowngradeHelper;
import com.codahale.metrics.Timer;
-import org.apache.hadoop.conf.Configuration;
import java.util.List;
import java.util.Map;
@@ -94,12 +93,12 @@ public class HoodieJavaWriteClient<T> extends
}
@Override
- protected HoodieTable createTable(HoodieWriteConfig config, Configuration
hadoopConf) {
+ protected HoodieTable createTable(HoodieWriteConfig config) {
return HoodieJavaTable.create(config, context);
}
@Override
- protected HoodieTable createTable(HoodieWriteConfig config, Configuration
hadoopConf, HoodieTableMetaClient metaClient) {
+ protected HoodieTable createTable(HoodieWriteConfig config,
HoodieTableMetaClient metaClient) {
return HoodieJavaTable.create(config, context, metaClient);
}
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java
index 6cf866482f7..f186c112e9b 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java
@@ -51,7 +51,6 @@ import org.apache.hudi.table.action.HoodieWriteMetadata;
import org.apache.hudi.table.upgrade.SparkUpgradeDowngradeHelper;
import com.codahale.metrics.Timer;
-import org.apache.hadoop.conf.Configuration;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.slf4j.Logger;
@@ -106,12 +105,12 @@ public class SparkRDDWriteClient<T> extends
}
@Override
- protected HoodieTable createTable(HoodieWriteConfig config, Configuration
hadoopConf) {
+ protected HoodieTable createTable(HoodieWriteConfig config) {
return HoodieSparkTable.create(config, context);
}
@Override
- protected HoodieTable createTable(HoodieWriteConfig config, Configuration
hadoopConf, HoodieTableMetaClient metaClient) {
+ protected HoodieTable createTable(HoodieWriteConfig config,
HoodieTableMetaClient metaClient) {
return HoodieSparkTable.create(config, context, metaClient);
}
diff --git
a/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/internal/DataSourceInternalWriterHelper.java
b/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/internal/DataSourceInternalWriterHelper.java
index 2193c9cb994..cb427aebedb 100644
---
a/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/internal/DataSourceInternalWriterHelper.java
+++
b/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/internal/DataSourceInternalWriterHelper.java
@@ -31,7 +31,6 @@ import org.apache.hudi.common.util.Option;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.exception.HoodieException;
import org.apache.hudi.storage.StorageConfiguration;
-import org.apache.hudi.table.HoodieSparkTable;
import org.apache.hudi.table.HoodieTable;
import org.apache.spark.api.java.JavaSparkContext;
@@ -67,12 +66,11 @@ public class DataSourceInternalWriterHelper {
this.writeClient = new SparkRDDWriteClient<>(new
HoodieSparkEngineContext(new JavaSparkContext(sparkSession.sparkContext())),
writeConfig);
this.writeClient.setOperationType(operationType);
this.writeClient.startCommitWithTime(instantTime);
- this.writeClient.initTable(operationType, Option.of(instantTime));
+ this.hoodieTable = this.writeClient.initTable(operationType,
Option.of(instantTime));
this.metaClient = HoodieTableMetaClient.builder()
.setConf(storageConf.newInstance()).setBasePath(writeConfig.getBasePath()).build();
this.metaClient.validateTableProperties(writeConfig.getProps());
- this.hoodieTable = HoodieSparkTable.create(writeConfig, new
HoodieSparkEngineContext(new JavaSparkContext(sparkSession.sparkContext())),
metaClient);
this.writeClient.preWrite(instantTime, WriteOperationType.BULK_INSERT,
metaClient);
}