This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 9189b05d3 [INLONG-7363][Sort] Fix iceberg connector null pointer error
(#7364)
9189b05d3 is described below
commit 9189b05d3b637989cb8bfa262398e72f7b3ecab0
Author: emhui <[email protected]>
AuthorDate: Tue Feb 14 15:57:49 2023 +0800
[INLONG-7363][Sort] Fix iceberg connector null pointer error (#7364)
---
.../src/main/java/org/apache/inlong/sort/iceberg/sink/FlinkSink.java | 5 +++--
1 file changed, 3 insertions(+), 2 deletions(-)
diff --git
a/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/FlinkSink.java
b/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/FlinkSink.java
index 2fc1ca13d..e816bc491 100644
---
a/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/FlinkSink.java
+++
b/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/FlinkSink.java
@@ -156,6 +156,7 @@ public class FlinkSink {
private Table table;
private TableSchema tableSchema;
private ActionsProvider actionProvider;
+ private boolean overwrite = false;
private boolean appendMode = false;
private Integer writeParallelism = null;
private List<String> equalityFieldColumns = null;
@@ -249,6 +250,7 @@ public class FlinkSink {
}
public Builder overwrite(boolean newOverwrite) {
+ this.overwrite = newOverwrite;
writeOptions.put(FlinkWriteOptions.OVERWRITE_MODE.key(),
Boolean.toString(newOverwrite));
return this;
}
@@ -525,8 +527,7 @@ public class FlinkSink {
private SingleOutputStreamOperator<Void> appendMultipleCommitter(
SingleOutputStreamOperator<MultipleWriteResult> writerStream) {
IcebergProcessOperator<MultipleWriteResult, Void>
multipleFilesCommiter =
- new IcebergProcessOperator<>(new
IcebergMultipleFilesCommiter(catalogLoader,
- flinkWriteConf.overwriteMode()));
+ new IcebergProcessOperator<>(new
IcebergMultipleFilesCommiter(catalogLoader, overwrite));
SingleOutputStreamOperator<Void> committerStream = writerStream
.transform(operatorName(ICEBERG_MULTIPLE_FILES_COMMITTER_NAME), Types.VOID,
multipleFilesCommiter)
.setParallelism(1)