diff --git a/core/src/main/java/org/apache/iceberg/io/BaseTaskWriter.java b/core/src/main/java/org/apache/iceberg/io/BaseTaskWriter.java index 0edf3662a18a..860a155deee6 100644 --- a/core/src/main/java/org/apache/iceberg/io/BaseTaskWriter.java +++ b/core/src/main/java/org/apache/iceberg/io/BaseTaskWriter.java @@ -56,6 +56,7 @@ public abstract class BaseTaskWriter implements TaskWriter { private final PartitionSpec spec; private final FileFormat format; private final FileAppenderFactory appenderFactory; + private final FileWriterFactory writerFactory; private final OutputFileFactory fileFactory; private final FileIO io; private final long targetFileSize; @@ -71,6 +72,23 @@ protected BaseTaskWriter( this.spec = spec; this.format = format; this.appenderFactory = appenderFactory; + this.writerFactory = null; + this.fileFactory = fileFactory; + this.io = io; + this.targetFileSize = targetFileSize; + } + + protected BaseTaskWriter( + PartitionSpec spec, + FileFormat format, + FileWriterFactory writerFactory, + OutputFileFactory fileFactory, + FileIO io, + long targetFileSize) { + this.spec = spec; + this.format = format; + this.appenderFactory = null; + this.writerFactory = writerFactory; this.fileFactory = fileFactory; this.io = io; this.targetFileSize = targetFileSize; @@ -138,7 +156,15 @@ protected BaseEqualityDeltaWriter( this.eqDeleteWriter = new RollingEqDeleteWriter(partition); this.posDeleteWriter = new SortingPositionOnlyDeleteWriter<>( - () -> appenderFactory.newPosDeleteWriter(newOutputFile(partition), format, partition), + () -> { + if (writerFactory != null) { + return writerFactory.newPositionDeleteWriter( + newOutputFile(partition), spec, partition); + } else { + return appenderFactory.newPosDeleteWriter( + newOutputFile(partition), format, partition); + } + }, deleteGranularity); this.insertedRowMap = StructLikeMap.create(deleteSchema.asStruct()); } @@ -388,7 +414,11 @@ public RollingFileWriter(StructLike partitionKey) { @Override DataWriter newWriter(EncryptedOutputFile file, StructLike partitionKey) { - return appenderFactory.newDataWriter(file, format, partitionKey); + if (writerFactory != null) { + return writerFactory.newDataWriter(file, spec, partitionKey); + } else { + return appenderFactory.newDataWriter(file, format, partitionKey); + } } @Override @@ -414,7 +444,11 @@ protected class RollingEqDeleteWriter extends BaseRollingWriter newWriter(EncryptedOutputFile file, StructLike partitionKey) { - return appenderFactory.newEqDeleteWriter(file, format, partitionKey); + if (writerFactory != null) { + return writerFactory.newEqualityDeleteWriter(file, spec, partitionKey); + } else { + return appenderFactory.newEqDeleteWriter(file, format, partitionKey); + } } @Override diff --git a/core/src/main/java/org/apache/iceberg/io/PartitionedFanoutWriter.java b/core/src/main/java/org/apache/iceberg/io/PartitionedFanoutWriter.java index b8c52024d2ed..8c196505ae14 100644 --- a/core/src/main/java/org/apache/iceberg/io/PartitionedFanoutWriter.java +++ b/core/src/main/java/org/apache/iceberg/io/PartitionedFanoutWriter.java @@ -38,6 +38,16 @@ protected PartitionedFanoutWriter( super(spec, format, appenderFactory, fileFactory, io, targetFileSize); } + protected PartitionedFanoutWriter( + PartitionSpec spec, + FileFormat format, + FileWriterFactory fileWriterFactory, + OutputFileFactory fileFactory, + FileIO io, + long targetFileSize) { + super(spec, format, fileWriterFactory, fileFactory, io, targetFileSize); + } + /** * Create a PartitionKey from the values in row. * diff --git a/core/src/main/java/org/apache/iceberg/io/PartitionedWriter.java b/core/src/main/java/org/apache/iceberg/io/PartitionedWriter.java index 625f8f94c997..50ca3dbb84e6 100644 --- a/core/src/main/java/org/apache/iceberg/io/PartitionedWriter.java +++ b/core/src/main/java/org/apache/iceberg/io/PartitionedWriter.java @@ -46,6 +46,16 @@ protected PartitionedWriter( super(spec, format, appenderFactory, fileFactory, io, targetFileSize); } + protected PartitionedWriter( + PartitionSpec spec, + FileFormat format, + FileWriterFactory fileWriterFactory, + OutputFileFactory fileFactory, + FileIO io, + long targetFileSize) { + super(spec, format, fileWriterFactory, fileFactory, io, targetFileSize); + } + /** * Create a PartitionKey from the values in row. * diff --git a/core/src/main/java/org/apache/iceberg/io/UnpartitionedWriter.java b/core/src/main/java/org/apache/iceberg/io/UnpartitionedWriter.java index 1c4aa3564bde..762b75a08fa2 100644 --- a/core/src/main/java/org/apache/iceberg/io/UnpartitionedWriter.java +++ b/core/src/main/java/org/apache/iceberg/io/UnpartitionedWriter.java @@ -37,6 +37,17 @@ public UnpartitionedWriter( currentWriter = new RollingFileWriter(null); } + public UnpartitionedWriter( + PartitionSpec spec, + FileFormat format, + FileWriterFactory fileWriterFactory, + OutputFileFactory fileFactory, + FileIO io, + long targetFileSize) { + super(spec, format, fileWriterFactory, fileFactory, io, targetFileSize); + currentWriter = new RollingFileWriter(null); + } + @Override public void write(T record) throws IOException { currentWriter.write(record); diff --git a/data/src/main/java/org/apache/iceberg/data/BaseFileWriterFactory.java b/data/src/main/java/org/apache/iceberg/data/BaseFileWriterFactory.java index 9e37c723be93..64aafbc8504f 100644 --- a/data/src/main/java/org/apache/iceberg/data/BaseFileWriterFactory.java +++ b/data/src/main/java/org/apache/iceberg/data/BaseFileWriterFactory.java @@ -19,6 +19,7 @@ package org.apache.iceberg.data; import java.io.IOException; +import java.io.Serializable; import java.io.UncheckedIOException; import java.util.Map; import org.apache.iceberg.FileFormat; @@ -37,9 +38,10 @@ import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.orc.ORC; import org.apache.iceberg.parquet.Parquet; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; /** A base writer factory to be extended by query engine integrations. */ -public abstract class BaseFileWriterFactory implements FileWriterFactory { +public abstract class BaseFileWriterFactory implements FileWriterFactory, Serializable { private final Table table; private final FileFormat dataFileFormat; private final Schema dataSchema; @@ -49,7 +51,32 @@ public abstract class BaseFileWriterFactory implements FileWriterFactory { private final Schema equalityDeleteRowSchema; private final SortOrder equalityDeleteSortOrder; private final Schema positionDeleteRowSchema; + private final Map writerProperties; + protected BaseFileWriterFactory( + Table table, + FileFormat dataFileFormat, + Schema dataSchema, + SortOrder dataSortOrder, + FileFormat deleteFileFormat, + int[] equalityFieldIds, + Schema equalityDeleteRowSchema, + SortOrder equalityDeleteSortOrder, + Schema positionDeleteRowSchema, + Map writerProperties) { + this.table = table; + this.dataFileFormat = dataFileFormat; + this.dataSchema = dataSchema; + this.dataSortOrder = dataSortOrder; + this.deleteFileFormat = deleteFileFormat; + this.equalityFieldIds = equalityFieldIds; + this.equalityDeleteRowSchema = equalityDeleteRowSchema; + this.equalityDeleteSortOrder = equalityDeleteSortOrder; + this.positionDeleteRowSchema = positionDeleteRowSchema; + this.writerProperties = writerProperties; + } + + @Deprecated protected BaseFileWriterFactory( Table table, FileFormat dataFileFormat, @@ -69,6 +96,7 @@ protected BaseFileWriterFactory( this.equalityDeleteRowSchema = equalityDeleteRowSchema; this.equalityDeleteSortOrder = equalityDeleteSortOrder; this.positionDeleteRowSchema = positionDeleteRowSchema; + this.writerProperties = ImmutableMap.of(); } protected abstract void configureDataWrite(Avro.DataWriteBuilder builder); @@ -103,6 +131,7 @@ public DataWriter newDataWriter( Avro.writeData(file) .schema(dataSchema) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .withSpec(spec) .withPartition(partition) @@ -119,6 +148,7 @@ public DataWriter newDataWriter( Parquet.writeData(file) .schema(dataSchema) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .withSpec(spec) .withPartition(partition) @@ -135,6 +165,7 @@ public DataWriter newDataWriter( ORC.writeData(file) .schema(dataSchema) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .withSpec(spec) .withPartition(partition) @@ -168,6 +199,7 @@ public EqualityDeleteWriter newEqualityDeleteWriter( Avro.DeleteWriteBuilder avroBuilder = Avro.writeDeletes(file) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .rowSchema(equalityDeleteRowSchema) .equalityFieldIds(equalityFieldIds) @@ -185,6 +217,7 @@ public EqualityDeleteWriter newEqualityDeleteWriter( Parquet.DeleteWriteBuilder parquetBuilder = Parquet.writeDeletes(file) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .rowSchema(equalityDeleteRowSchema) .equalityFieldIds(equalityFieldIds) @@ -202,6 +235,7 @@ public EqualityDeleteWriter newEqualityDeleteWriter( ORC.DeleteWriteBuilder orcBuilder = ORC.writeDeletes(file) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .rowSchema(equalityDeleteRowSchema) .equalityFieldIds(equalityFieldIds) @@ -237,6 +271,7 @@ public PositionDeleteWriter newPositionDeleteWriter( Avro.DeleteWriteBuilder avroBuilder = Avro.writeDeletes(file) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .rowSchema(positionDeleteRowSchema) .withSpec(spec) @@ -252,6 +287,7 @@ public PositionDeleteWriter newPositionDeleteWriter( Parquet.DeleteWriteBuilder parquetBuilder = Parquet.writeDeletes(file) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .rowSchema(positionDeleteRowSchema) .withSpec(spec) @@ -267,6 +303,7 @@ public PositionDeleteWriter newPositionDeleteWriter( ORC.DeleteWriteBuilder orcBuilder = ORC.writeDeletes(file) .setAll(properties) + .setAll(writerProperties) .metricsConfig(metricsConfig) .rowSchema(positionDeleteRowSchema) .withSpec(spec) diff --git a/data/src/test/java/org/apache/iceberg/io/TestFileWriterFactory.java b/data/src/test/java/org/apache/iceberg/io/TestFileWriterFactory.java index ab1d295125f2..8354149bf495 100644 --- a/data/src/test/java/org/apache/iceberg/io/TestFileWriterFactory.java +++ b/data/src/test/java/org/apache/iceberg/io/TestFileWriterFactory.java @@ -22,6 +22,7 @@ import static org.apache.iceberg.MetadataColumns.DELETE_FILE_POS; import static org.apache.iceberg.MetadataColumns.DELETE_FILE_ROW_FIELD_NAME; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatNoException; import static org.assertj.core.api.Assumptions.assumeThat; import java.io.File; @@ -55,8 +56,10 @@ import org.apache.iceberg.types.Types; import org.apache.iceberg.util.CharSequenceSet; import org.apache.iceberg.util.Pair; +import org.apache.iceberg.util.SerializationUtil; import org.apache.iceberg.util.StructLikeSet; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; import org.junit.jupiter.api.TestTemplate; import org.junit.jupiter.api.extension.ExtendWith; @@ -406,6 +409,15 @@ public void testPositionDeleteWriterMultipleDataFiles() throws IOException { assertThat(actualRowSet("*")).isEqualTo(toSet(expectedRows)); } + @Test + void testSerialization() { + FileWriterFactory writerFactory = newWriterFactory(table.schema()); + assertThatNoException().isThrownBy(() -> SerializationUtil.serializeToBytes(writerFactory)); + + byte[] serialized = SerializationUtil.serializeToBytes(writerFactory); + assertThatNoException().isThrownBy(() -> SerializationUtil.deserializeFromBytes(serialized)); + } + private DataFile writeData( FileWriterFactory writerFactory, List rows, PartitionSpec spec, StructLike partitionKey) throws IOException { diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/BaseDeltaTaskWriter.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/BaseDeltaTaskWriter.java index d845046cd2f6..f68eff691266 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/BaseDeltaTaskWriter.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/BaseDeltaTaskWriter.java @@ -32,8 +32,8 @@ import org.apache.iceberg.flink.RowDataWrapper; import org.apache.iceberg.flink.data.RowDataProjection; import org.apache.iceberg.io.BaseTaskWriter; -import org.apache.iceberg.io.FileAppenderFactory; import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.OutputFileFactory; import org.apache.iceberg.relocated.com.google.common.collect.Sets; import org.apache.iceberg.types.TypeUtil; @@ -50,7 +50,7 @@ abstract class BaseDeltaTaskWriter extends BaseTaskWriter { BaseDeltaTaskWriter( PartitionSpec spec, FileFormat format, - FileAppenderFactory appenderFactory, + FileWriterFactory fileWriterFactory, OutputFileFactory fileFactory, FileIO io, long targetFileSize, @@ -58,7 +58,7 @@ abstract class BaseDeltaTaskWriter extends BaseTaskWriter { RowType flinkSchema, Set equalityFieldIds, boolean upsert) { - super(spec, format, appenderFactory, fileFactory, io, targetFileSize); + super(spec, format, fileWriterFactory, fileFactory, io, targetFileSize); this.schema = schema; this.deleteSchema = TypeUtil.select(schema, Sets.newHashSet(equalityFieldIds)); this.wrapper = new RowDataWrapper(flinkSchema, schema.asStruct()); diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkAppenderFactory.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkAppenderFactory.java index b6f1392d1562..d21bbec81fed 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkAppenderFactory.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkAppenderFactory.java @@ -48,6 +48,11 @@ import org.apache.iceberg.parquet.Parquet; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +/** + * @deprecated Deprecated as of 1.11.0 in favor of {@link FlinkFileWriterFactory}. This class will + * be removed in the 1.12.0. + */ +@Deprecated public class FlinkAppenderFactory implements FileAppenderFactory, Serializable { private final Schema schema; private final RowType flinkSchema; diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkFileWriterFactory.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkFileWriterFactory.java index 2183fe062af4..948dd252d5de 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkFileWriterFactory.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkFileWriterFactory.java @@ -42,13 +42,14 @@ import org.apache.iceberg.orc.ORC; import org.apache.iceberg.parquet.Parquet; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -class FlinkFileWriterFactory extends BaseFileWriterFactory implements Serializable { +public class FlinkFileWriterFactory extends BaseFileWriterFactory implements Serializable { private RowType dataFlinkType; private RowType equalityDeleteFlinkType; private RowType positionDeleteFlinkType; - FlinkFileWriterFactory( + private FlinkFileWriterFactory( Table table, FileFormat dataFileFormat, Schema dataSchema, @@ -60,7 +61,8 @@ class FlinkFileWriterFactory extends BaseFileWriterFactory implements S RowType equalityDeleteFlinkType, SortOrder equalityDeleteSortOrder, Schema positionDeleteRowSchema, - RowType positionDeleteFlinkType) { + RowType positionDeleteFlinkType, + Map writeProperties) { super( table, @@ -71,7 +73,8 @@ class FlinkFileWriterFactory extends BaseFileWriterFactory implements S equalityFieldIds, equalityDeleteRowSchema, equalityDeleteSortOrder, - positionDeleteRowSchema); + positionDeleteRowSchema, + writeProperties); this.dataFlinkType = dataFlinkType; this.equalityDeleteFlinkType = equalityDeleteFlinkType; @@ -170,7 +173,7 @@ private RowType positionDeleteFlinkType() { return positionDeleteFlinkType; } - static class Builder { + public static class Builder { private final Table table; private FileFormat dataFileFormat; private Schema dataSchema; @@ -183,8 +186,9 @@ static class Builder { private SortOrder equalityDeleteSortOrder; private Schema positionDeleteRowSchema; private RowType positionDeleteFlinkType; + private Map writerProperties = ImmutableMap.of(); - Builder(Table table) { + public Builder(Table table) { this.table = table; Map properties = table.properties(); @@ -198,12 +202,12 @@ static class Builder { this.deleteFileFormat = FileFormat.fromString(deleteFileFormatName); } - Builder dataFileFormat(FileFormat newDataFileFormat) { + public Builder dataFileFormat(FileFormat newDataFileFormat) { this.dataFileFormat = newDataFileFormat; return this; } - Builder dataSchema(Schema newDataSchema) { + public Builder dataSchema(Schema newDataSchema) { this.dataSchema = newDataSchema; return this; } @@ -213,27 +217,27 @@ Builder dataSchema(Schema newDataSchema) { * *

If not set, the value is derived from the provided Iceberg schema. */ - Builder dataFlinkType(RowType newDataFlinkType) { + public Builder dataFlinkType(RowType newDataFlinkType) { this.dataFlinkType = newDataFlinkType; return this; } - Builder dataSortOrder(SortOrder newDataSortOrder) { + public Builder dataSortOrder(SortOrder newDataSortOrder) { this.dataSortOrder = newDataSortOrder; return this; } - Builder deleteFileFormat(FileFormat newDeleteFileFormat) { + public Builder deleteFileFormat(FileFormat newDeleteFileFormat) { this.deleteFileFormat = newDeleteFileFormat; return this; } - Builder equalityFieldIds(int[] newEqualityFieldIds) { + public Builder equalityFieldIds(int[] newEqualityFieldIds) { this.equalityFieldIds = newEqualityFieldIds; return this; } - Builder equalityDeleteRowSchema(Schema newEqualityDeleteRowSchema) { + public Builder equalityDeleteRowSchema(Schema newEqualityDeleteRowSchema) { this.equalityDeleteRowSchema = newEqualityDeleteRowSchema; return this; } @@ -243,17 +247,17 @@ Builder equalityDeleteRowSchema(Schema newEqualityDeleteRowSchema) { * *

If not set, the value is derived from the provided Iceberg schema. */ - Builder equalityDeleteFlinkType(RowType newEqualityDeleteFlinkType) { + public Builder equalityDeleteFlinkType(RowType newEqualityDeleteFlinkType) { this.equalityDeleteFlinkType = newEqualityDeleteFlinkType; return this; } - Builder equalityDeleteSortOrder(SortOrder newEqualityDeleteSortOrder) { + public Builder equalityDeleteSortOrder(SortOrder newEqualityDeleteSortOrder) { this.equalityDeleteSortOrder = newEqualityDeleteSortOrder; return this; } - Builder positionDeleteRowSchema(Schema newPositionDeleteRowSchema) { + public Builder positionDeleteRowSchema(Schema newPositionDeleteRowSchema) { this.positionDeleteRowSchema = newPositionDeleteRowSchema; return this; } @@ -263,12 +267,18 @@ Builder positionDeleteRowSchema(Schema newPositionDeleteRowSchema) { * *

If not set, the value is derived from the provided Iceberg schema. */ - Builder positionDeleteFlinkType(RowType newPositionDeleteFlinkType) { + public Builder positionDeleteFlinkType(RowType newPositionDeleteFlinkType) { this.positionDeleteFlinkType = newPositionDeleteFlinkType; return this; } - FlinkFileWriterFactory build() { + /** Sets default writer properties. */ + public Builder writerProperties(Map newWriterProperties) { + this.writerProperties = newWriterProperties; + return this; + } + + public FlinkFileWriterFactory build() { boolean noEqualityDeleteConf = equalityFieldIds == null && equalityDeleteRowSchema == null; boolean fullEqualityDeleteConf = equalityFieldIds != null && equalityDeleteRowSchema != null; Preconditions.checkArgument( @@ -287,7 +297,8 @@ FlinkFileWriterFactory build() { equalityDeleteFlinkType, equalityDeleteSortOrder, positionDeleteRowSchema, - positionDeleteFlinkType); + positionDeleteFlinkType, + writerProperties); } } } diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/PartitionedDeltaWriter.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/PartitionedDeltaWriter.java index 3eb4dba80281..afbc14b7f153 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/PartitionedDeltaWriter.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/PartitionedDeltaWriter.java @@ -28,8 +28,8 @@ import org.apache.iceberg.PartitionKey; import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.Schema; -import org.apache.iceberg.io.FileAppenderFactory; import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.OutputFileFactory; import org.apache.iceberg.relocated.com.google.common.collect.Maps; import org.apache.iceberg.util.Tasks; @@ -43,7 +43,7 @@ class PartitionedDeltaWriter extends BaseDeltaTaskWriter { PartitionedDeltaWriter( PartitionSpec spec, FileFormat format, - FileAppenderFactory appenderFactory, + FileWriterFactory fileWriterFactory, OutputFileFactory fileFactory, FileIO io, long targetFileSize, @@ -54,7 +54,7 @@ class PartitionedDeltaWriter extends BaseDeltaTaskWriter { super( spec, format, - appenderFactory, + fileWriterFactory, fileFactory, io, targetFileSize, diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/RowDataTaskWriterFactory.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/RowDataTaskWriterFactory.java index 7c11b20c449d..ef2c795e2322 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/RowDataTaskWriterFactory.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/RowDataTaskWriterFactory.java @@ -30,8 +30,8 @@ import org.apache.iceberg.Schema; import org.apache.iceberg.Table; import org.apache.iceberg.flink.RowDataWrapper; -import org.apache.iceberg.io.FileAppenderFactory; import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.OutputFileFactory; import org.apache.iceberg.io.PartitionedFanoutWriter; import org.apache.iceberg.io.TaskWriter; @@ -51,7 +51,7 @@ public class RowDataTaskWriterFactory implements TaskWriterFactory { private final FileFormat format; private final Set equalityFieldIds; private final boolean upsert; - private final FileAppenderFactory appenderFactory; + private final FileWriterFactory fileWriterFactory; private transient OutputFileFactory outputFileFactory; @@ -122,36 +122,40 @@ public RowDataTaskWriterFactory( this.upsert = upsert; if (equalityFieldIds == null || equalityFieldIds.isEmpty()) { - this.appenderFactory = - new FlinkAppenderFactory( - table, schema, flinkSchema, writeProperties, spec, null, null, null); + this.fileWriterFactory = + new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(format) + .dataSchema(schema) + .dataFlinkType(flinkSchema) + .writerProperties(writeProperties) + .build(); } else if (upsert) { // In upsert mode, only the new row is emitted using INSERT row kind. Therefore, any column of // the inserted row // may differ from the deleted row other than the primary key fields, and the delete file must // contain values // that are correct for the deleted row. Therefore, only write the equality delete fields. - this.appenderFactory = - new FlinkAppenderFactory( - table, - schema, - flinkSchema, - writeProperties, - spec, - ArrayUtil.toPrimitive(equalityFieldIds.toArray(new Integer[0])), - TypeUtil.select(schema, Sets.newHashSet(equalityFieldIds)), - null); + this.fileWriterFactory = + new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(format) + .dataSchema(schema) + .dataFlinkType(flinkSchema) + .deleteFileFormat(format) + .equalityFieldIds(ArrayUtil.toPrimitive(equalityFieldIds.toArray(new Integer[0]))) + .equalityDeleteRowSchema(TypeUtil.select(schema, Sets.newHashSet(equalityFieldIds))) + .writerProperties(writeProperties) + .build(); } else { - this.appenderFactory = - new FlinkAppenderFactory( - table, - schema, - flinkSchema, - writeProperties, - spec, - ArrayUtil.toPrimitive(equalityFieldIds.toArray(new Integer[0])), - schema, - null); + this.fileWriterFactory = + new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(format) + .dataSchema(schema) + .dataFlinkType(flinkSchema) + .deleteFileFormat(format) + .equalityFieldIds(ArrayUtil.toPrimitive(equalityFieldIds.toArray(new Integer[0]))) + .equalityDeleteRowSchema(schema) + .writerProperties(writeProperties) + .build(); } } @@ -189,7 +193,7 @@ public TaskWriter create() { return new UnpartitionedWriter<>( spec, format, - appenderFactory, + fileWriterFactory, outputFileFactory, tableSupplier.get().io(), targetFileSizeBytes); @@ -197,7 +201,7 @@ public TaskWriter create() { return new RowDataPartitionedFanoutWriter( spec, format, - appenderFactory, + fileWriterFactory, outputFileFactory, tableSupplier.get().io(), targetFileSizeBytes, @@ -210,7 +214,7 @@ public TaskWriter create() { return new UnpartitionedDeltaWriter( spec, format, - appenderFactory, + fileWriterFactory, outputFileFactory, tableSupplier.get().io(), targetFileSizeBytes, @@ -222,7 +226,7 @@ public TaskWriter create() { return new PartitionedDeltaWriter( spec, format, - appenderFactory, + fileWriterFactory, outputFileFactory, tableSupplier.get().io(), targetFileSizeBytes, @@ -248,13 +252,13 @@ private static class RowDataPartitionedFanoutWriter extends PartitionedFanoutWri RowDataPartitionedFanoutWriter( PartitionSpec spec, FileFormat format, - FileAppenderFactory appenderFactory, + FileWriterFactory fileWriterFactory, OutputFileFactory fileFactory, FileIO io, long targetFileSize, Schema schema, RowType flinkSchema) { - super(spec, format, appenderFactory, fileFactory, io, targetFileSize); + super(spec, format, fileWriterFactory, fileFactory, io, targetFileSize); this.partitionKey = new PartitionKey(spec, schema); this.rowDataWrapper = new RowDataWrapper(flinkSchema, schema.asStruct()); } diff --git a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/UnpartitionedDeltaWriter.java b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/UnpartitionedDeltaWriter.java index b6ad03514bb0..e709206c9499 100644 --- a/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/UnpartitionedDeltaWriter.java +++ b/flink/v2.0/flink/src/main/java/org/apache/iceberg/flink/sink/UnpartitionedDeltaWriter.java @@ -25,8 +25,8 @@ import org.apache.iceberg.FileFormat; import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.Schema; -import org.apache.iceberg.io.FileAppenderFactory; import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.OutputFileFactory; class UnpartitionedDeltaWriter extends BaseDeltaTaskWriter { @@ -35,7 +35,7 @@ class UnpartitionedDeltaWriter extends BaseDeltaTaskWriter { UnpartitionedDeltaWriter( PartitionSpec spec, FileFormat format, - FileAppenderFactory appenderFactory, + FileWriterFactory fileWriterFactory, OutputFileFactory fileFactory, FileIO io, long targetFileSize, @@ -46,7 +46,7 @@ class UnpartitionedDeltaWriter extends BaseDeltaTaskWriter { super( spec, format, - appenderFactory, + fileWriterFactory, fileFactory, io, targetFileSize, diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/SimpleDataUtil.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/SimpleDataUtil.java index d9c9f7ad3f02..6d7a9bfd6b84 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/SimpleDataUtil.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/SimpleDataUtil.java @@ -38,7 +38,6 @@ import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.iceberg.DataFile; -import org.apache.iceberg.DataFiles; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileFormat; import org.apache.iceberg.FileScanTask; @@ -56,16 +55,18 @@ import org.apache.iceberg.deletes.EqualityDeleteWriter; import org.apache.iceberg.deletes.PositionDelete; import org.apache.iceberg.deletes.PositionDeleteWriter; +import org.apache.iceberg.encryption.EncryptedFiles; import org.apache.iceberg.encryption.EncryptedOutputFile; -import org.apache.iceberg.flink.sink.FlinkAppenderFactory; -import org.apache.iceberg.hadoop.HadoopInputFile; +import org.apache.iceberg.encryption.EncryptionKeyMetadata; +import org.apache.iceberg.flink.sink.FlinkFileWriterFactory; import org.apache.iceberg.hadoop.HadoopTables; import org.apache.iceberg.io.CloseableIterable; -import org.apache.iceberg.io.FileAppender; +import org.apache.iceberg.io.DataWriter; import org.apache.iceberg.io.FileAppenderFactory; +import org.apache.iceberg.io.FileWriterFactory; +import org.apache.iceberg.io.OutputFile; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; import org.apache.iceberg.relocated.com.google.common.collect.Iterables; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Maps; @@ -175,26 +176,19 @@ public static DataFile writeFile( FileFormat fileFormat = FileFormat.fromFileName(filename); Preconditions.checkNotNull(fileFormat, "Cannot determine format for file: %s", filename); - RowType flinkSchema = FlinkSchemaUtil.convert(schema); - FileAppenderFactory appenderFactory = - new FlinkAppenderFactory( - table, schema, flinkSchema, ImmutableMap.of(), spec, null, null, null); + FileWriterFactory writerFactory = + new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(fileFormat) + .dataSchema(schema) + .build(); - FileAppender appender = appenderFactory.newAppender(fromPath(path, conf), fileFormat); - try (FileAppender closeableAppender = appender) { - closeableAppender.addAll(rows); - } - - DataFiles.Builder builder = - DataFiles.builder(spec) - .withInputFile(HadoopInputFile.fromPath(path, conf)) - .withMetrics(appender.metrics()); + DataWriter writer = + writerFactory.newDataWriter(encrypt(fromPath(path, conf)), spec, partition); - if (partition != null) { - builder = builder.withPartition(partition); - } + writer.write(rows); + writer.close(); - return builder.build(); + return writer.toDataFile(); } public static DeleteFile writeEqDeleteFile( @@ -205,9 +199,7 @@ public static DeleteFile writeEqDeleteFile( List deletes) throws IOException { EncryptedOutputFile outputFile = - table - .encryption() - .encrypt(fromPath(new Path(table.location(), filename), new Configuration())); + encrypt(fromPath(new Path(table.location(), filename), new Configuration())); EqualityDeleteWriter eqWriter = appenderFactory.newEqDeleteWriter(outputFile, format, null); @@ -217,6 +209,25 @@ public static DeleteFile writeEqDeleteFile( return eqWriter.toDeleteFile(); } + public static DeleteFile writeEqDeleteFile( + Table table, + PartitionSpec spec, + String filename, + FileWriterFactory writerFactory, + List deletes) + throws IOException { + EncryptedOutputFile outputFile = + encrypt(fromPath(new Path(table.location(), filename), new Configuration())); + + EqualityDeleteWriter eqWriter = + writerFactory.newEqualityDeleteWriter(outputFile, spec, null); + try (EqualityDeleteWriter writer = eqWriter) { + writer.write(deletes); + } + + return eqWriter.toDeleteFile(); + } + public static DeleteFile writePosDeleteFile( Table table, FileFormat format, @@ -225,9 +236,7 @@ public static DeleteFile writePosDeleteFile( List> positions) throws IOException { EncryptedOutputFile outputFile = - table - .encryption() - .encrypt(fromPath(new Path(table.location(), filename), new Configuration())); + encrypt(fromPath(new Path(table.location(), filename), new Configuration())); PositionDeleteWriter posWriter = appenderFactory.newPosDeleteWriter(outputFile, format, null); @@ -240,6 +249,28 @@ public static DeleteFile writePosDeleteFile( return posWriter.toDeleteFile(); } + public static DeleteFile writePosDeleteFile( + Table table, + PartitionSpec spec, + String filename, + FileWriterFactory writerFactory, + List> positions) + throws IOException { + EncryptedOutputFile outputFile = + encrypt(fromPath(new Path(table.location(), filename), new Configuration())); + + PositionDeleteWriter posWriter = + writerFactory.newPositionDeleteWriter(outputFile, spec, null); + PositionDelete posDelete = PositionDelete.create(); + try (PositionDeleteWriter writer = posWriter) { + for (Pair p : positions) { + writer.write(posDelete.set(p.first(), p.second())); + } + } + + return posWriter.toDeleteFile(); + } + private static List convertToRecords(List rows) { List records = Lists.newArrayList(); for (RowData row : rows) { @@ -466,4 +497,8 @@ public static List matchingPartitions( }) .collect(Collectors.toList()); } + + private static EncryptedOutputFile encrypt(OutputFile out) { + return EncryptedFiles.encryptedOutput(out, EncryptionKeyMetadata.EMPTY); + } } diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestCompressionSettings.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestCompressionSettings.java index 5a74db5713a5..da5b5f6c28f0 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestCompressionSettings.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestCompressionSettings.java @@ -35,10 +35,12 @@ import org.apache.iceberg.Table; import org.apache.iceberg.TableProperties; import org.apache.iceberg.common.DynFields; +import org.apache.iceberg.data.BaseFileWriterFactory; import org.apache.iceberg.flink.FlinkWriteConf; import org.apache.iceberg.flink.FlinkWriteOptions; import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.io.BaseTaskWriter; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.TaskWriter; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; import org.junit.jupiter.api.BeforeEach; @@ -244,14 +246,14 @@ private static Map appenderProperties( DynFields.builder() .hiddenImpl(IcebergStreamWriter.class, "writer") .build(operatorField.get()); - DynFields.BoundField appenderField = + DynFields.BoundField writerFactoryField = DynFields.builder() - .hiddenImpl(BaseTaskWriter.class, "appenderFactory") + .hiddenImpl(BaseTaskWriter.class, "writerFactory") .build(writerField.get()); DynFields.BoundField> propsField = DynFields.builder() - .hiddenImpl(FlinkAppenderFactory.class, "props") - .build(appenderField.get()); + .hiddenImpl(BaseFileWriterFactory.class, "writerProperties") + .build(writerFactoryField.get()); return propsField.get(); } } diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkManifest.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkManifest.java index c21c3d5cc21b..10197ddfaff6 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkManifest.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkManifest.java @@ -39,10 +39,8 @@ import org.apache.iceberg.ManifestFile; import org.apache.iceberg.ManifestFiles; import org.apache.iceberg.Table; -import org.apache.iceberg.flink.FlinkSchemaUtil; import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TestHelpers; -import org.apache.iceberg.io.FileAppenderFactory; import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; import org.apache.iceberg.relocated.com.google.common.collect.Lists; @@ -57,7 +55,7 @@ public class TestFlinkManifest { @TempDir protected Path temporaryFolder; private Table table; - private FileAppenderFactory appenderFactory; + private FlinkFileWriterFactory fileWriterFactory; private final AtomicInteger fileCount = new AtomicInteger(0); @BeforeEach @@ -75,16 +73,15 @@ public void before() throws IOException { new int[] { table.schema().findField("id").fieldId(), table.schema().findField("data").fieldId() }; - this.appenderFactory = - new FlinkAppenderFactory( - table, - table.schema(), - FlinkSchemaUtil.convert(table.schema()), - table.properties(), - table.spec(), - equalityFieldIds, - table.schema(), - null); + this.fileWriterFactory = + new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(FileFormat.PARQUET) + .dataSchema(table.schema()) + .deleteFileFormat(FileFormat.PARQUET) + .equalityFieldIds(equalityFieldIds) + .equalityDeleteRowSchema(table.schema()) + .writerProperties(table.properties()) + .build(); } @Test @@ -262,13 +259,13 @@ private DataFile writeDataFile(String filename, List rows) throws IOExc private DeleteFile writeEqDeleteFile(String filename, List deletes) throws IOException { return SimpleDataUtil.writeEqDeleteFile( - table, FileFormat.PARQUET, filename, appenderFactory, deletes); + table, table.spec(), filename, fileWriterFactory, deletes); } private DeleteFile writePosDeleteFile(String filename, List> positions) throws IOException { return SimpleDataUtil.writePosDeleteFile( - table, FileFormat.PARQUET, filename, appenderFactory, positions); + table, table.spec(), filename, fileWriterFactory, positions); } private List generateDataFiles(int fileNum) throws IOException { diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java index 76338a185a62..2fb7fc10a8a1 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergCommitter.java @@ -80,11 +80,10 @@ import org.apache.iceberg.Table; import org.apache.iceberg.TableProperties; import org.apache.iceberg.TestBase; -import org.apache.iceberg.flink.FlinkSchemaUtil; import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.flink.TestHelpers; -import org.apache.iceberg.io.FileAppenderFactory; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; @@ -1070,7 +1069,7 @@ public void testDeleteFiles() throws Exception { assumeThat(formatVersion).as("Only support delete in format v2").isGreaterThanOrEqualTo(2); - FileAppenderFactory appenderFactory = createDeletableAppenderFactory(); + FileWriterFactory writerFactory = createFileWriterFactory(); try (OneInputStreamOperatorTestHarness< CommittableMessage, CommittableMessage> @@ -1125,7 +1124,8 @@ public void testDeleteFiles() throws Exception { checkpointId = 3; RowData delete1 = SimpleDataUtil.createDelete(1, "aaa"); DeleteFile deleteFile1 = - writeEqDeleteFile(appenderFactory, "delete-file-1", ImmutableList.of(delete1)); + writeEqDeleteFile( + writerFactory, "delete-file-1", ImmutableList.of(delete1), table.spec()); RowData row4 = SimpleDataUtil.createInsert(4, "ddd"); DataFile dataFile4 = writeDataFile("data-file-4", ImmutableList.of(row4)); @@ -1235,20 +1235,18 @@ private CommittableSummary processElement( return processElement(withRecord, myJobID, checkpointId, testHarness, subTaskId, operatorId); } - private FileAppenderFactory createDeletableAppenderFactory() { + private FileWriterFactory createFileWriterFactory() { int[] equalityFieldIds = new int[] { table.schema().findField("id").fieldId(), table.schema().findField("data").fieldId() }; - return new FlinkAppenderFactory( - table, - table.schema(), - FlinkSchemaUtil.convert(table.schema()), - table.properties(), - table.spec(), - equalityFieldIds, - table.schema(), - null); + return new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(FileFormat.PARQUET) + .dataSchema(table.schema()) + .deleteFileFormat(FileFormat.PARQUET) + .equalityFieldIds(equalityFieldIds) + .equalityDeleteRowSchema(table.schema()) + .build(); } private List assertFlinkManifests(int expectedCount) throws IOException { @@ -1272,10 +1270,12 @@ private DataFile writeDataFile(String filename, List rows) throws IOExc } private DeleteFile writeEqDeleteFile( - FileAppenderFactory appenderFactory, String filename, List deletes) + FileWriterFactory writerFactory, + String filename, + List deletes, + PartitionSpec spec) throws IOException { - return SimpleDataUtil.writeEqDeleteFile( - table, FileFormat.PARQUET, filename, appenderFactory, deletes); + return SimpleDataUtil.writeEqDeleteFile(table, spec, filename, writerFactory, deletes); } private OneInputStreamOperatorTestHarness< diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java index 65fb9b8f69b4..68621018be57 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/TestIcebergFilesCommitter.java @@ -65,11 +65,11 @@ import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.StructLike; import org.apache.iceberg.TestBase; -import org.apache.iceberg.flink.FlinkSchemaUtil; import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TestHelpers; import org.apache.iceberg.flink.TestTableLoader; import org.apache.iceberg.io.FileAppenderFactory; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; import org.apache.iceberg.relocated.com.google.common.collect.Lists; @@ -786,7 +786,7 @@ public void testDeleteFiles() throws Exception { JobID jobId = new JobID(); OperatorID operatorId; - FileAppenderFactory appenderFactory = createDeletableAppenderFactory(); + FileWriterFactory writerFactory = createWriterFactory(); try (OneInputStreamOperatorTestHarness harness = createStreamSink(jobId)) { @@ -829,7 +829,7 @@ public void testDeleteFiles() throws Exception { RowData delete1 = SimpleDataUtil.createDelete(1, "aaa"); DeleteFile deleteFile1 = - writeEqDeleteFile(appenderFactory, "delete-file-1", ImmutableList.of(delete1)); + writeEqDeleteFile(writerFactory, "delete-file-1", ImmutableList.of(delete1)); assertMaxCommittedCheckpointId(jobId, operatorId, checkpoint); harness.processElement( new FlinkWriteResult( @@ -860,7 +860,7 @@ public void testCommitTwoCheckpointsInSingleTxn() throws Exception { JobID jobId = new JobID(); OperatorID operatorId; - FileAppenderFactory appenderFactory = createDeletableAppenderFactory(); + FileWriterFactory writerFactory = createWriterFactory(); try (OneInputStreamOperatorTestHarness harness = createStreamSink(jobId)) { @@ -875,7 +875,7 @@ public void testCommitTwoCheckpointsInSingleTxn() throws Exception { RowData delete3 = SimpleDataUtil.createDelete(3, "ccc"); DataFile dataFile1 = writeDataFile("data-file-1", ImmutableList.of(insert1, insert2)); DeleteFile deleteFile1 = - writeEqDeleteFile(appenderFactory, "delete-file-1", ImmutableList.of(delete3)); + writeEqDeleteFile(writerFactory, "delete-file-1", ImmutableList.of(delete3)); harness.processElement( new FlinkWriteResult( checkpoint, @@ -889,7 +889,7 @@ public void testCommitTwoCheckpointsInSingleTxn() throws Exception { RowData delete2 = SimpleDataUtil.createDelete(2, "bbb"); DataFile dataFile2 = writeDataFile("data-file-2", ImmutableList.of(insert4)); DeleteFile deleteFile2 = - writeEqDeleteFile(appenderFactory, "delete-file-2", ImmutableList.of(delete2)); + writeEqDeleteFile(writerFactory, "delete-file-2", ImmutableList.of(delete2)); harness.processElement( new FlinkWriteResult( ++checkpoint, @@ -937,16 +937,14 @@ public void testCommitMultipleCheckpointsForV2Table() throws Exception { JobID jobId = new JobID(); OperatorID operatorId; - FileAppenderFactory appenderFactory = - new FlinkAppenderFactory( - table, - table.schema(), - FlinkSchemaUtil.convert(table.schema()), - table.properties(), - table.spec(), - new int[] {table.schema().findField("id").fieldId()}, - table.schema(), - null); + FileWriterFactory writerFactory = + new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(format) + .dataSchema(table.schema()) + .deleteFileFormat(format) + .equalityFieldIds(new int[] {table.schema().findField("id").fieldId()}) + .equalityDeleteRowSchema(table.schema()) + .build(); try (OneInputStreamOperatorTestHarness harness = createStreamSink(jobId)) { @@ -964,7 +962,7 @@ public void testCommitMultipleCheckpointsForV2Table() throws Exception { DataFile dataFile = writeDataFile("data-file-" + i, ImmutableList.of(insert1, insert2)); DeleteFile deleteFile = writeEqDeleteFile( - appenderFactory, "delete-file-" + i, ImmutableList.of(insert1, insert2)); + writerFactory, "delete-file-" + i, ImmutableList.of(insert1, insert2)); harness.processElement( new FlinkWriteResult( ++checkpoint, @@ -1090,9 +1088,9 @@ private int getStagingManifestSpecId(OperatorStateStore operatorStateStore, long } private DeleteFile writeEqDeleteFile( - FileAppenderFactory appenderFactory, String filename, List deletes) + FileWriterFactory writerFactory, String filename, List deletes) throws IOException { - return SimpleDataUtil.writeEqDeleteFile(table, format, filename, appenderFactory, deletes); + return SimpleDataUtil.writeEqDeleteFile(table, table.spec(), filename, writerFactory, deletes); } private DeleteFile writePosDeleteFile( @@ -1103,20 +1101,18 @@ private DeleteFile writePosDeleteFile( return SimpleDataUtil.writePosDeleteFile(table, format, filename, appenderFactory, positions); } - private FileAppenderFactory createDeletableAppenderFactory() { + private FileWriterFactory createWriterFactory() { int[] equalityFieldIds = new int[] { table.schema().findField("id").fieldId(), table.schema().findField("data").fieldId() }; - return new FlinkAppenderFactory( - table, - table.schema(), - FlinkSchemaUtil.convert(table.schema()), - table.properties(), - table.spec(), - equalityFieldIds, - table.schema(), - null); + return new FlinkFileWriterFactory.Builder(table) + .dataFileFormat(format) + .dataSchema(table.schema()) + .deleteFileFormat(format) + .equalityFieldIds(equalityFieldIds) + .equalityDeleteRowSchema(table.schema()) + .build(); } private ManifestFile createTestingManifestFile(Path manifestPath, DataFile dataFile) diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestDynamicWriter.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestDynamicWriter.java index 91a3a5d5a7c6..89dd4f22594a 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestDynamicWriter.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/sink/dynamic/TestDynamicWriter.java @@ -33,10 +33,11 @@ import org.apache.iceberg.catalog.Catalog; import org.apache.iceberg.catalog.TableIdentifier; import org.apache.iceberg.common.DynFields; +import org.apache.iceberg.data.BaseFileWriterFactory; import org.apache.iceberg.flink.SimpleDataUtil; -import org.apache.iceberg.flink.sink.FlinkAppenderFactory; import org.apache.iceberg.flink.sink.TestFlinkIcebergSinkBase; import org.apache.iceberg.io.BaseTaskWriter; +import org.apache.iceberg.io.FileWriterFactory; import org.apache.iceberg.io.TaskWriter; import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; @@ -241,14 +242,14 @@ private Map properties(DynamicWriter dynamicWriter) { DynFields.BoundField>> writerField = DynFields.builder().hiddenImpl(dynamicWriter.getClass(), "writers").build(dynamicWriter); - DynFields.BoundField appenderField = + DynFields.BoundField writerFactoryField = DynFields.builder() - .hiddenImpl(BaseTaskWriter.class, "appenderFactory") + .hiddenImpl(BaseTaskWriter.class, "writerFactory") .build(writerField.get().values().iterator().next()); DynFields.BoundField> propsField = DynFields.builder() - .hiddenImpl(FlinkAppenderFactory.class, "props") - .build(appenderField.get()); + .hiddenImpl(BaseFileWriterFactory.class, "writerProperties") + .build(writerFactoryField.get()); return propsField.get(); } }