diff --git a/connectors/connector-hive/hive-bridge/src/main/java/org/apache/flink/connectors/hive/InputOutputFormat.java b/connectors/connector-hive/hive-bridge/src/main/java/org/apache/flink/connectors/hive/InputOutputFormat.java index 072c64c92..b0c4e37e7 100644 --- a/connectors/connector-hive/hive-bridge/src/main/java/org/apache/flink/connectors/hive/InputOutputFormat.java +++ b/connectors/connector-hive/hive-bridge/src/main/java/org/apache/flink/connectors/hive/InputOutputFormat.java @@ -54,7 +54,7 @@ public class InputOutputFormat { } public static OutputFormat createOutputFormat( - Catalog catalog, DynamicTableFactory.Context context, Map partitions) { + Catalog catalog, DynamicTableFactory.Context context, Map partitions, Boolean overwriteSink) { if (!(catalog instanceof HiveCatalog)) { throw new RuntimeException("Catalog should be hive catalog."); @@ -72,6 +72,8 @@ public static OutputFormat createOutputFormat( hiveTableSink.applyStaticPartition(partitions); } + hiveTableSink.applyOverwrite(overwriteSink); + return hiveTableSink.getOutputFormat(); } } diff --git a/core/src/main/java/com/alibaba/alink/common/io/catalog/HiveCatalog.java b/core/src/main/java/com/alibaba/alink/common/io/catalog/HiveCatalog.java index d69760764..7b66dbd05 100644 --- a/core/src/main/java/com/alibaba/alink/common/io/catalog/HiveCatalog.java +++ b/core/src/main/java/com/alibaba/alink/common/io/catalog/HiveCatalog.java @@ -833,6 +833,7 @@ public boolean isTemporary() { return factory.doAsThrowRuntime(() -> { String partitionSpec = params.get(HiveCatalogParams.PARTITION); + boolean overwriteSink = params.get(HiveCatalogParams.OVERWRITE_SINK); Map partitions = null; if (!StringUtils.isNullOrWhitespaceOnly(partitionSpec)) { @@ -844,10 +845,10 @@ public boolean isTemporary() { true, Thread.currentThread().getContextClassLoader() ); - Method method = inputOutputFormat.getMethod("createOutputFormat", Catalog.class, DynamicTableFactory.Context.class, Map.class); + Method method = inputOutputFormat.getMethod("createOutputFormat", Catalog.class, DynamicTableFactory.Context.class, Map.class, Boolean.class); OutputFormat internalRet = - (OutputFormat ) method.invoke(null, catalog, context, partitions); + (OutputFormat ) method.invoke(null, catalog, context, partitions, overwriteSink); return new RichOutputFormatWithClassLoader(factory, internalRet); });