1
0

Revert "[HUDI-3951]support generan parameter 'sink.parallelism' for flink-hudi (#5405)" (#5421)

This reverts commit bda3db078e.
This commit is contained in:
ForwardXu
2022-04-25 12:58:27 +08:00
committed by GitHub
parent d994c58cc0
commit 9054b85961
2 changed files with 1 additions and 5 deletions

View File

@@ -34,7 +34,6 @@ import org.apache.hudi.keygen.constant.KeyGeneratorType;
import org.apache.flink.configuration.ConfigOption; import org.apache.flink.configuration.ConfigOption;
import org.apache.flink.configuration.ConfigOptions; import org.apache.flink.configuration.ConfigOptions;
import org.apache.flink.configuration.Configuration; import org.apache.flink.configuration.Configuration;
import org.apache.flink.table.factories.FactoryUtil;
import java.lang.reflect.Field; import java.lang.reflect.Field;
import java.util.ArrayList; import java.util.ArrayList;
@@ -233,8 +232,6 @@ public class FlinkOptions extends HoodieConfig {
// ------------------------------------------------------------------------ // ------------------------------------------------------------------------
// Write Options // Write Options
// ------------------------------------------------------------------------ // ------------------------------------------------------------------------
public static final ConfigOption<Integer> SINK_PARALLELISM = FactoryUtil.SINK_PARALLELISM;
public static final ConfigOption<String> TABLE_NAME = ConfigOptions public static final ConfigOption<String> TABLE_NAME = ConfigOptions
.key(HoodieWriteConfig.TBL_NAME.key()) .key(HoodieWriteConfig.TBL_NAME.key())
.stringType() .stringType()

View File

@@ -82,8 +82,7 @@ public class HoodieTableSink implements DynamicTableSink, SupportsPartitioning,
} }
// default parallelism // default parallelism
int parallelism = conf.getInteger(FlinkOptions.SINK_PARALLELISM, int parallelism = dataStream.getExecutionConfig().getParallelism();
dataStream.getExecutionConfig().getParallelism());
DataStream<Object> pipeline; DataStream<Object> pipeline;
// bootstrap // bootstrap
final DataStream<HoodieRecord> hoodieRecordDataStream = final DataStream<HoodieRecord> hoodieRecordDataStream =