From a6a25fccaa4ecd67283f93e4e5e619589a3f8905 Mon Sep 17 00:00:00 2001 From: Steven Wu Date: Tue, 16 Mar 2021 22:09:41 -0700 Subject: [PATCH 1/5] Flink: refactor DataIterator to use composition (instead of inheritance) --- .../iceberg/flink/source/DataIterator.java | 52 ++++----------- .../flink/source/DecryptedInputFiles.java | 64 +++++++++++++++++++ .../flink/source/FlinkInputFormat.java | 11 ++-- .../iceberg/flink/source/IteratorReader.java | 33 ++++++++++ ...erator.java => RowDataIteratorReader.java} | 54 ++++++++-------- .../iceberg/flink/source/RowDataRewriter.java | 6 +- 6 files changed, 147 insertions(+), 73 deletions(-) create mode 100644 flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java create mode 100644 flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java rename flink/src/main/java/org/apache/iceberg/flink/source/{RowDataIterator.java => RowDataIteratorReader.java} (71%) diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java index f74a8968fab8..4ab8c3908565 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java @@ -21,64 +21,35 @@ import java.io.IOException; import java.io.UncheckedIOException; -import java.nio.ByteBuffer; import java.util.Iterator; -import java.util.Map; -import java.util.stream.Stream; import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.FileScanTask; -import org.apache.iceberg.encryption.EncryptedFiles; -import org.apache.iceberg.encryption.EncryptedInputFile; import org.apache.iceberg.encryption.EncryptionManager; import org.apache.iceberg.io.CloseableIterator; import org.apache.iceberg.io.FileIO; -import org.apache.iceberg.io.InputFile; -import org.apache.iceberg.relocated.com.google.common.base.Preconditions; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.iceberg.relocated.com.google.common.collect.Maps; /** * Base class of Flink iterators. * * @param is the Java class returned by this iterator whose objects contain one or more rows. */ -abstract class DataIterator implements CloseableIterator { +public class DataIterator implements CloseableIterator { - private Iterator tasks; - private final Map inputFiles; + private final IteratorReader iteratorReader; + private final DecryptedInputFiles decryptedInputFiles; + private Iterator tasks; private CloseableIterator currentIterator; - DataIterator(CombinedScanTask task, FileIO io, EncryptionManager encryption) { - this.tasks = task.files().iterator(); - - Map keyMetadata = Maps.newHashMap(); - task.files().stream() - .flatMap(fileScanTask -> Stream.concat(Stream.of(fileScanTask.file()), fileScanTask.deletes().stream())) - .forEach(file -> keyMetadata.put(file.path().toString(), file.keyMetadata())); - Stream encrypted = keyMetadata.entrySet().stream() - .map(entry -> EncryptedFiles.encryptedInput(io.newInputFile(entry.getKey()), entry.getValue())); - - // decrypt with the batch call to avoid multiple RPCs to a key server, if possible - Iterable decryptedFiles = encryption.decrypt(encrypted::iterator); - - Map files = Maps.newHashMapWithExpectedSize(task.files().size()); - decryptedFiles.forEach(decrypted -> files.putIfAbsent(decrypted.location(), decrypted)); - this.inputFiles = ImmutableMap.copyOf(files); + public DataIterator(IteratorReader iteratorReader, CombinedScanTask task, + FileIO io, EncryptionManager encryption) { + this.iteratorReader = iteratorReader; + this.decryptedInputFiles = new DecryptedInputFiles(task, io, encryption); + this.tasks = task.files().iterator(); this.currentIterator = CloseableIterator.empty(); } - InputFile getInputFile(FileScanTask task) { - Preconditions.checkArgument(!task.isDataTask(), "Invalid task type"); - - return inputFiles.get(task.file().path().toString()); - } - - InputFile getInputFile(String location) { - return inputFiles.get(location); - } - @Override public boolean hasNext() { updateCurrentIterator(); @@ -106,7 +77,9 @@ private void updateCurrentIterator() { } } - abstract CloseableIterator openTaskIterator(FileScanTask scanTask) throws IOException; + private CloseableIterator openTaskIterator(FileScanTask scanTask) throws IOException { + return iteratorReader.open(scanTask, decryptedInputFiles); + } @Override public void close() throws IOException { @@ -114,4 +87,5 @@ public void close() throws IOException { currentIterator.close(); tasks = null; } + } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java b/flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java new file mode 100644 index 000000000000..99461d1dc9ef --- /dev/null +++ b/flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java @@ -0,0 +1,64 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iceberg.flink.source; + +import java.nio.ByteBuffer; +import java.util.Collections; +import java.util.Map; +import java.util.stream.Stream; +import org.apache.iceberg.CombinedScanTask; +import org.apache.iceberg.FileScanTask; +import org.apache.iceberg.encryption.EncryptedFiles; +import org.apache.iceberg.encryption.EncryptedInputFile; +import org.apache.iceberg.encryption.EncryptionManager; +import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.InputFile; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.Maps; + +class DecryptedInputFiles { + + private final Map inputFiles; + + DecryptedInputFiles(CombinedScanTask combinedTask, FileIO io, EncryptionManager encryption) { + Map keyMetadata = Maps.newHashMap(); + combinedTask.files().stream() + .flatMap(fileScanTask -> Stream.concat(Stream.of(fileScanTask.file()), fileScanTask.deletes().stream())) + .forEach(file -> keyMetadata.put(file.path().toString(), file.keyMetadata())); + Stream encrypted = keyMetadata.entrySet().stream() + .map(entry -> EncryptedFiles.encryptedInput(io.newInputFile(entry.getKey()), entry.getValue())); + + // decrypt with the batch call to avoid multiple RPCs to a key server, if possible + Iterable decryptedFiles = encryption.decrypt(encrypted::iterator); + + Map files = Maps.newHashMapWithExpectedSize(combinedTask.files().size()); + decryptedFiles.forEach(decrypted -> files.putIfAbsent(decrypted.location(), decrypted)); + this.inputFiles = Collections.unmodifiableMap(files); + } + + InputFile getInputFile(FileScanTask task) { + Preconditions.checkArgument(!task.isDataTask(), "Invalid task type"); + return inputFiles.get(task.file().path().toString()); + } + + InputFile getInputFile(String location) { + return inputFiles.get(location); + } +} diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java b/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java index 1bad1c25952e..f0bb2b0c32c4 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java @@ -42,21 +42,22 @@ public class FlinkInputFormat extends RichInputFormat private static final long serialVersionUID = 1L; private final TableLoader tableLoader; - private final Schema tableSchema; private final FileIO io; private final EncryptionManager encryption; private final ScanContext context; + private final RowDataIteratorReader rowDataReader; - private transient RowDataIterator iterator; + private transient DataIterator iterator; private transient long currentReadCount = 0L; FlinkInputFormat(TableLoader tableLoader, Schema tableSchema, FileIO io, EncryptionManager encryption, ScanContext context) { this.tableLoader = tableLoader; - this.tableSchema = tableSchema; this.io = io; this.encryption = encryption; this.context = context; + this.rowDataReader = new RowDataIteratorReader(tableSchema, + context.project(), context.nameMapping(), context.caseSensitive()); } @VisibleForTesting @@ -91,9 +92,7 @@ public void configure(Configuration parameters) { @Override public void open(FlinkInputSplit split) { - this.iterator = new RowDataIterator( - split.getTask(), io, encryption, tableSchema, context.project(), context.nameMapping(), - context.caseSensitive()); + this.iterator = new DataIterator<>(rowDataReader, split.getTask(), io, encryption); } @Override diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java b/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java new file mode 100644 index 000000000000..6d129048d667 --- /dev/null +++ b/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java @@ -0,0 +1,33 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iceberg.flink.source; + +import java.io.Serializable; +import org.apache.iceberg.FileScanTask; +import org.apache.iceberg.io.CloseableIterator; + +/** + * Read a {@link FileScanTask} into a {@link CloseableIterator} + */ +public interface IteratorReader extends Serializable { + + CloseableIterator open(FileScanTask fileScanTask, DecryptedInputFiles decryptedInputFiles); + +} diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIterator.java b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java similarity index 71% rename from flink/src/main/java/org/apache/iceberg/flink/source/RowDataIterator.java rename to flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java index 5a568144d1f7..00da4102e2d6 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIterator.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java @@ -21,14 +21,12 @@ import java.util.Map; import org.apache.flink.table.data.RowData; -import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.Schema; import org.apache.iceberg.StructLike; import org.apache.iceberg.avro.Avro; import org.apache.iceberg.data.DeleteFilter; -import org.apache.iceberg.encryption.EncryptionManager; import org.apache.iceberg.flink.FlinkSchemaUtil; import org.apache.iceberg.flink.RowDataWrapper; import org.apache.iceberg.flink.data.FlinkAvroReader; @@ -37,7 +35,6 @@ import org.apache.iceberg.flink.data.RowDataUtil; import org.apache.iceberg.io.CloseableIterable; import org.apache.iceberg.io.CloseableIterator; -import org.apache.iceberg.io.FileIO; import org.apache.iceberg.io.InputFile; import org.apache.iceberg.mapping.NameMappingParser; import org.apache.iceberg.orc.ORC; @@ -47,16 +44,15 @@ import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.util.PartitionUtil; -class RowDataIterator extends DataIterator { +public class RowDataIteratorReader implements IteratorReader { private final Schema tableSchema; private final Schema projectedSchema; private final String nameMapping; private final boolean caseSensitive; - RowDataIterator(CombinedScanTask task, FileIO io, EncryptionManager encryption, Schema tableSchema, - Schema projectedSchema, String nameMapping, boolean caseSensitive) { - super(task, io, encryption); + public RowDataIteratorReader(Schema tableSchema, Schema projectedSchema, + String nameMapping, boolean caseSensitive) { this.tableSchema = tableSchema; this.projectedSchema = projectedSchema; this.nameMapping = nameMapping; @@ -64,34 +60,35 @@ class RowDataIterator extends DataIterator { } @Override - protected CloseableIterator openTaskIterator(FileScanTask task) { + public CloseableIterator open(FileScanTask task, DecryptedInputFiles decryptedInputFiles) { Schema partitionSchema = TypeUtil.select(projectedSchema, task.spec().identitySourceIds()); Map idToConstant = partitionSchema.columns().isEmpty() ? ImmutableMap.of() : PartitionUtil.constantsMap(task, RowDataUtil::convertConstant); - FlinkDeleteFilter deletes = new FlinkDeleteFilter(task, tableSchema, projectedSchema); - CloseableIterable iterable = deletes.filter(newIterable(task, deletes.requiredSchema(), idToConstant)); - - return iterable.iterator(); + FlinkDeleteFilter deletes = new FlinkDeleteFilter(task, tableSchema, projectedSchema, decryptedInputFiles); + return deletes + .filter(newIterable(task, deletes.requiredSchema(), idToConstant, decryptedInputFiles)) + .iterator(); } - private CloseableIterable newIterable(FileScanTask task, Schema schema, Map idToConstant) { + private CloseableIterable newIterable( + FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { CloseableIterable iter; if (task.isDataTask()) { throw new UnsupportedOperationException("Cannot read data task."); } else { switch (task.file().format()) { case PARQUET: - iter = newParquetIterable(task, schema, idToConstant); + iter = newParquetIterable(task, schema, idToConstant, decryptedInputFiles); break; case AVRO: - iter = newAvroIterable(task, schema, idToConstant); + iter = newAvroIterable(task, schema, idToConstant, decryptedInputFiles); break; case ORC: - iter = newOrcIterable(task, schema, idToConstant); + iter = newOrcIterable(task, schema, idToConstant, decryptedInputFiles); break; default: @@ -103,8 +100,9 @@ private CloseableIterable newIterable(FileScanTask task, Schema schema, return iter; } - private CloseableIterable newAvroIterable(FileScanTask task, Schema schema, Map idToConstant) { - Avro.ReadBuilder builder = Avro.read(getInputFile(task)) + private CloseableIterable newAvroIterable( + FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { + Avro.ReadBuilder builder = Avro.read(decryptedInputFiles.getInputFile(task)) .reuseContainers() .project(schema) .split(task.start(), task.length()) @@ -117,9 +115,9 @@ private CloseableIterable newAvroIterable(FileScanTask task, Schema sch return builder.build(); } - private CloseableIterable newParquetIterable(FileScanTask task, Schema schema, - Map idToConstant) { - Parquet.ReadBuilder builder = Parquet.read(getInputFile(task)) + private CloseableIterable newParquetIterable( + FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { + Parquet.ReadBuilder builder = Parquet.read(decryptedInputFiles.getInputFile(task)) .reuseContainers() .split(task.start(), task.length()) .project(schema) @@ -135,11 +133,12 @@ private CloseableIterable newParquetIterable(FileScanTask task, Schema return builder.build(); } - private CloseableIterable newOrcIterable(FileScanTask task, Schema schema, Map idToConstant) { + private CloseableIterable newOrcIterable( + FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { Schema readSchemaWithoutConstantAndMetadataFields = TypeUtil.selectNot(schema, Sets.union(idToConstant.keySet(), MetadataColumns.metadataFieldIds())); - ORC.ReadBuilder builder = ORC.read(getInputFile(task)) + ORC.ReadBuilder builder = ORC.read(decryptedInputFiles.getInputFile(task)) .project(readSchemaWithoutConstantAndMetadataFields) .split(task.start(), task.length()) .createReaderFunc(readOrcSchema -> new FlinkOrcReader(schema, readOrcSchema, idToConstant)) @@ -153,12 +152,15 @@ private CloseableIterable newOrcIterable(FileScanTask task, Schema sche return builder.build(); } - private class FlinkDeleteFilter extends DeleteFilter { + private static class FlinkDeleteFilter extends DeleteFilter { private final RowDataWrapper asStructLike; + private final DecryptedInputFiles decryptedInputFiles; - FlinkDeleteFilter(FileScanTask task, Schema tableSchema, Schema requestedSchema) { + FlinkDeleteFilter(FileScanTask task, Schema tableSchema, Schema requestedSchema, + DecryptedInputFiles decryptedInputFiles) { super(task, tableSchema, requestedSchema); this.asStructLike = new RowDataWrapper(FlinkSchemaUtil.convert(requiredSchema()), requiredSchema().asStruct()); + this.decryptedInputFiles = decryptedInputFiles; } @Override @@ -168,7 +170,7 @@ protected StructLike asStructLike(RowData row) { @Override protected InputFile getInputFile(String location) { - return RowDataIterator.this.getInputFile(location); + return decryptedInputFiles.getInputFile(location); } } } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java index a6cd374c3044..d21e4a5d1517 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java @@ -99,6 +99,7 @@ public static class RewriteMap extends RichMapFunction taskWriterFactory; + private final RowDataIteratorReader rowDataReader; public RewriteMap(Schema schema, String nameMapping, FileIO io, boolean caseSensitive, EncryptionManager encryptionManager, TaskWriterFactory taskWriterFactory) { @@ -108,6 +109,7 @@ public RewriteMap(Schema schema, String nameMapping, FileIO io, boolean caseSens this.caseSensitive = caseSensitive; this.encryptionManager = encryptionManager; this.taskWriterFactory = taskWriterFactory; + this.rowDataReader = new RowDataIteratorReader(schema, schema, nameMapping, caseSensitive); } @Override @@ -122,8 +124,8 @@ public void open(Configuration parameters) { public List map(CombinedScanTask task) throws Exception { // Initialize the task writer. this.writer = taskWriterFactory.create(); - try (RowDataIterator iterator = - new RowDataIterator(task, io, encryptionManager, schema, schema, nameMapping, caseSensitive)) { + try (DataIterator iterator = + new DataIterator<>(rowDataReader, task, io, encryptionManager)) { while (iterator.hasNext()) { RowData rowData = iterator.next(); writer.write(rowData); From 2b123ca19dcbd0f9c62d1bf6fcfb62327096b742 Mon Sep 17 00:00:00 2001 From: Steven Wu Date: Fri, 13 Aug 2021 13:25:44 -0700 Subject: [PATCH 2/5] rename DecryptedInputFiles to InputFilesDecryptor and move it to the core module also removed an unnecessary whitespace change --- .../encryption/InputFilesDecryptor.java | 56 ++++++++++++++++ .../iceberg/flink/source/DataIterator.java | 8 +-- .../flink/source/DecryptedInputFiles.java | 64 ------------------- .../iceberg/flink/source/IteratorReader.java | 3 +- .../flink/source/RowDataIteratorReader.java | 35 +++++----- 5 files changed, 80 insertions(+), 86 deletions(-) create mode 100644 core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java delete mode 100644 flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java diff --git a/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java b/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java new file mode 100644 index 000000000000..54c48c23764d --- /dev/null +++ b/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java @@ -0,0 +1,56 @@ +/* + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.iceberg.encryption; + +import java.nio.ByteBuffer; +import java.util.Collections; +import java.util.Map; +import java.util.stream.Stream; +import org.apache.iceberg.CombinedScanTask; +import org.apache.iceberg.FileScanTask; +import org.apache.iceberg.io.FileIO; +import org.apache.iceberg.io.InputFile; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.Maps; + +public class InputFilesDecryptor { + + private final Map decryptedInputFiles; + + public InputFilesDecryptor(CombinedScanTask combinedTask, FileIO io, EncryptionManager encryption) { + Map keyMetadata = Maps.newHashMap(); + combinedTask.files().stream() + .flatMap(fileScanTask -> Stream.concat(Stream.of(fileScanTask.file()), fileScanTask.deletes().stream())) + .forEach(file -> keyMetadata.put(file.path().toString(), file.keyMetadata())); + Stream encrypted = keyMetadata.entrySet().stream() + .map(entry -> EncryptedFiles.encryptedInput(io.newInputFile(entry.getKey()), entry.getValue())); + + // decrypt with the batch call to avoid multiple RPCs to a key server, if possible + Iterable decryptedFiles = encryption.decrypt(encrypted::iterator); + + Map files = Maps.newHashMapWithExpectedSize(combinedTask.files().size()); + decryptedFiles.forEach(decrypted -> files.putIfAbsent(decrypted.location(), decrypted)); + this.decryptedInputFiles = Collections.unmodifiableMap(files); + } + + public InputFile getInputFile(FileScanTask task) { + Preconditions.checkArgument(!task.isDataTask(), "Invalid task type"); + return decryptedInputFiles.get(task.file().path().toString()); + } + + public InputFile getInputFile(String location) { + return decryptedInputFiles.get(location); + } +} diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java index 4ab8c3908565..ff3cdf594dbf 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java @@ -25,6 +25,7 @@ import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.encryption.EncryptionManager; +import org.apache.iceberg.encryption.InputFilesDecryptor; import org.apache.iceberg.io.CloseableIterator; import org.apache.iceberg.io.FileIO; @@ -37,7 +38,7 @@ public class DataIterator implements CloseableIterator { private final IteratorReader iteratorReader; - private final DecryptedInputFiles decryptedInputFiles; + private final InputFilesDecryptor inputFilesDecryptor; private Iterator tasks; private CloseableIterator currentIterator; @@ -45,7 +46,7 @@ public DataIterator(IteratorReader iteratorReader, CombinedScanTask task, FileIO io, EncryptionManager encryption) { this.iteratorReader = iteratorReader; - this.decryptedInputFiles = new DecryptedInputFiles(task, io, encryption); + this.inputFilesDecryptor = new InputFilesDecryptor(task, io, encryption); this.tasks = task.files().iterator(); this.currentIterator = CloseableIterator.empty(); } @@ -78,7 +79,7 @@ private void updateCurrentIterator() { } private CloseableIterator openTaskIterator(FileScanTask scanTask) throws IOException { - return iteratorReader.open(scanTask, decryptedInputFiles); + return iteratorReader.open(scanTask, inputFilesDecryptor); } @Override @@ -87,5 +88,4 @@ public void close() throws IOException { currentIterator.close(); tasks = null; } - } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java b/flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java deleted file mode 100644 index 99461d1dc9ef..000000000000 --- a/flink/src/main/java/org/apache/iceberg/flink/source/DecryptedInputFiles.java +++ /dev/null @@ -1,64 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ - -package org.apache.iceberg.flink.source; - -import java.nio.ByteBuffer; -import java.util.Collections; -import java.util.Map; -import java.util.stream.Stream; -import org.apache.iceberg.CombinedScanTask; -import org.apache.iceberg.FileScanTask; -import org.apache.iceberg.encryption.EncryptedFiles; -import org.apache.iceberg.encryption.EncryptedInputFile; -import org.apache.iceberg.encryption.EncryptionManager; -import org.apache.iceberg.io.FileIO; -import org.apache.iceberg.io.InputFile; -import org.apache.iceberg.relocated.com.google.common.base.Preconditions; -import org.apache.iceberg.relocated.com.google.common.collect.Maps; - -class DecryptedInputFiles { - - private final Map inputFiles; - - DecryptedInputFiles(CombinedScanTask combinedTask, FileIO io, EncryptionManager encryption) { - Map keyMetadata = Maps.newHashMap(); - combinedTask.files().stream() - .flatMap(fileScanTask -> Stream.concat(Stream.of(fileScanTask.file()), fileScanTask.deletes().stream())) - .forEach(file -> keyMetadata.put(file.path().toString(), file.keyMetadata())); - Stream encrypted = keyMetadata.entrySet().stream() - .map(entry -> EncryptedFiles.encryptedInput(io.newInputFile(entry.getKey()), entry.getValue())); - - // decrypt with the batch call to avoid multiple RPCs to a key server, if possible - Iterable decryptedFiles = encryption.decrypt(encrypted::iterator); - - Map files = Maps.newHashMapWithExpectedSize(combinedTask.files().size()); - decryptedFiles.forEach(decrypted -> files.putIfAbsent(decrypted.location(), decrypted)); - this.inputFiles = Collections.unmodifiableMap(files); - } - - InputFile getInputFile(FileScanTask task) { - Preconditions.checkArgument(!task.isDataTask(), "Invalid task type"); - return inputFiles.get(task.file().path().toString()); - } - - InputFile getInputFile(String location) { - return inputFiles.get(location); - } -} diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java b/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java index 6d129048d667..2e1e87e501b9 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java @@ -21,6 +21,7 @@ import java.io.Serializable; import org.apache.iceberg.FileScanTask; +import org.apache.iceberg.encryption.InputFilesDecryptor; import org.apache.iceberg.io.CloseableIterator; /** @@ -28,6 +29,6 @@ */ public interface IteratorReader extends Serializable { - CloseableIterator open(FileScanTask fileScanTask, DecryptedInputFiles decryptedInputFiles); + CloseableIterator open(FileScanTask fileScanTask, InputFilesDecryptor inputFilesDecryptor); } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java index 00da4102e2d6..36668bde6b14 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java @@ -27,6 +27,7 @@ import org.apache.iceberg.StructLike; import org.apache.iceberg.avro.Avro; import org.apache.iceberg.data.DeleteFilter; +import org.apache.iceberg.encryption.InputFilesDecryptor; import org.apache.iceberg.flink.FlinkSchemaUtil; import org.apache.iceberg.flink.RowDataWrapper; import org.apache.iceberg.flink.data.FlinkAvroReader; @@ -60,35 +61,35 @@ public RowDataIteratorReader(Schema tableSchema, Schema projectedSchema, } @Override - public CloseableIterator open(FileScanTask task, DecryptedInputFiles decryptedInputFiles) { + public CloseableIterator open(FileScanTask task, InputFilesDecryptor inputFilesDecryptor) { Schema partitionSchema = TypeUtil.select(projectedSchema, task.spec().identitySourceIds()); Map idToConstant = partitionSchema.columns().isEmpty() ? ImmutableMap.of() : PartitionUtil.constantsMap(task, RowDataUtil::convertConstant); - FlinkDeleteFilter deletes = new FlinkDeleteFilter(task, tableSchema, projectedSchema, decryptedInputFiles); + FlinkDeleteFilter deletes = new FlinkDeleteFilter(task, tableSchema, projectedSchema, inputFilesDecryptor); return deletes - .filter(newIterable(task, deletes.requiredSchema(), idToConstant, decryptedInputFiles)) + .filter(newIterable(task, deletes.requiredSchema(), idToConstant, inputFilesDecryptor)) .iterator(); } private CloseableIterable newIterable( - FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { + FileScanTask task, Schema schema, Map idToConstant, InputFilesDecryptor inputFilesDecryptor) { CloseableIterable iter; if (task.isDataTask()) { throw new UnsupportedOperationException("Cannot read data task."); } else { switch (task.file().format()) { case PARQUET: - iter = newParquetIterable(task, schema, idToConstant, decryptedInputFiles); + iter = newParquetIterable(task, schema, idToConstant, inputFilesDecryptor); break; case AVRO: - iter = newAvroIterable(task, schema, idToConstant, decryptedInputFiles); + iter = newAvroIterable(task, schema, idToConstant, inputFilesDecryptor); break; case ORC: - iter = newOrcIterable(task, schema, idToConstant, decryptedInputFiles); + iter = newOrcIterable(task, schema, idToConstant, inputFilesDecryptor); break; default: @@ -101,8 +102,8 @@ private CloseableIterable newIterable( } private CloseableIterable newAvroIterable( - FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { - Avro.ReadBuilder builder = Avro.read(decryptedInputFiles.getInputFile(task)) + FileScanTask task, Schema schema, Map idToConstant, InputFilesDecryptor inputFilesDecryptor) { + Avro.ReadBuilder builder = Avro.read(inputFilesDecryptor.getInputFile(task)) .reuseContainers() .project(schema) .split(task.start(), task.length()) @@ -116,8 +117,8 @@ private CloseableIterable newAvroIterable( } private CloseableIterable newParquetIterable( - FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { - Parquet.ReadBuilder builder = Parquet.read(decryptedInputFiles.getInputFile(task)) + FileScanTask task, Schema schema, Map idToConstant, InputFilesDecryptor inputFilesDecryptor) { + Parquet.ReadBuilder builder = Parquet.read(inputFilesDecryptor.getInputFile(task)) .reuseContainers() .split(task.start(), task.length()) .project(schema) @@ -134,11 +135,11 @@ private CloseableIterable newParquetIterable( } private CloseableIterable newOrcIterable( - FileScanTask task, Schema schema, Map idToConstant, DecryptedInputFiles decryptedInputFiles) { + FileScanTask task, Schema schema, Map idToConstant, InputFilesDecryptor inputFilesDecryptor) { Schema readSchemaWithoutConstantAndMetadataFields = TypeUtil.selectNot(schema, Sets.union(idToConstant.keySet(), MetadataColumns.metadataFieldIds())); - ORC.ReadBuilder builder = ORC.read(decryptedInputFiles.getInputFile(task)) + ORC.ReadBuilder builder = ORC.read(inputFilesDecryptor.getInputFile(task)) .project(readSchemaWithoutConstantAndMetadataFields) .split(task.start(), task.length()) .createReaderFunc(readOrcSchema -> new FlinkOrcReader(schema, readOrcSchema, idToConstant)) @@ -154,13 +155,13 @@ private CloseableIterable newOrcIterable( private static class FlinkDeleteFilter extends DeleteFilter { private final RowDataWrapper asStructLike; - private final DecryptedInputFiles decryptedInputFiles; + private final InputFilesDecryptor inputFilesDecryptor; FlinkDeleteFilter(FileScanTask task, Schema tableSchema, Schema requestedSchema, - DecryptedInputFiles decryptedInputFiles) { + InputFilesDecryptor inputFilesDecryptor) { super(task, tableSchema, requestedSchema); this.asStructLike = new RowDataWrapper(FlinkSchemaUtil.convert(requiredSchema()), requiredSchema().asStruct()); - this.decryptedInputFiles = decryptedInputFiles; + this.inputFilesDecryptor = inputFilesDecryptor; } @Override @@ -170,7 +171,7 @@ protected StructLike asStructLike(RowData row) { @Override protected InputFile getInputFile(String location) { - return decryptedInputFiles.getInputFile(location); + return inputFilesDecryptor.getInputFile(location); } } } From ced59a9b94a8c9f525933ed92717f1f64d2e62a1 Mon Sep 17 00:00:00 2001 From: Steven Wu Date: Fri, 13 Aug 2021 13:30:52 -0700 Subject: [PATCH 3/5] rename IteratorReader to FileReader --- .../org/apache/iceberg/flink/source/DataIterator.java | 10 +++++----- .../source/{IteratorReader.java => FileReader.java} | 2 +- .../apache/iceberg/flink/source/FlinkInputFormat.java | 4 ++-- ...wDataIteratorReader.java => RowDataFileReader.java} | 6 +++--- .../apache/iceberg/flink/source/RowDataRewriter.java | 4 ++-- 5 files changed, 13 insertions(+), 13 deletions(-) rename flink/src/main/java/org/apache/iceberg/flink/source/{IteratorReader.java => FileReader.java} (95%) rename flink/src/main/java/org/apache/iceberg/flink/source/{RowDataIteratorReader.java => RowDataFileReader.java} (96%) diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java index ff3cdf594dbf..79a44ee884cc 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java @@ -36,15 +36,15 @@ */ public class DataIterator implements CloseableIterator { - private final IteratorReader iteratorReader; + private final FileReader fileReader; private final InputFilesDecryptor inputFilesDecryptor; private Iterator tasks; private CloseableIterator currentIterator; - public DataIterator(IteratorReader iteratorReader, CombinedScanTask task, - FileIO io, EncryptionManager encryption) { - this.iteratorReader = iteratorReader; + public DataIterator(FileReader fileReader, CombinedScanTask task, + FileIO io, EncryptionManager encryption) { + this.fileReader = fileReader; this.inputFilesDecryptor = new InputFilesDecryptor(task, io, encryption); this.tasks = task.files().iterator(); @@ -79,7 +79,7 @@ private void updateCurrentIterator() { } private CloseableIterator openTaskIterator(FileScanTask scanTask) throws IOException { - return iteratorReader.open(scanTask, inputFilesDecryptor); + return fileReader.open(scanTask, inputFilesDecryptor); } @Override diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java b/flink/src/main/java/org/apache/iceberg/flink/source/FileReader.java similarity index 95% rename from flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java rename to flink/src/main/java/org/apache/iceberg/flink/source/FileReader.java index 2e1e87e501b9..2e339cb4067f 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/IteratorReader.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/FileReader.java @@ -27,7 +27,7 @@ /** * Read a {@link FileScanTask} into a {@link CloseableIterator} */ -public interface IteratorReader extends Serializable { +public interface FileReader extends Serializable { CloseableIterator open(FileScanTask fileScanTask, InputFilesDecryptor inputFilesDecryptor); diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java b/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java index f0bb2b0c32c4..d614c3ccaf17 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java @@ -45,7 +45,7 @@ public class FlinkInputFormat extends RichInputFormat private final FileIO io; private final EncryptionManager encryption; private final ScanContext context; - private final RowDataIteratorReader rowDataReader; + private final RowDataFileReader rowDataReader; private transient DataIterator iterator; private transient long currentReadCount = 0L; @@ -56,7 +56,7 @@ public class FlinkInputFormat extends RichInputFormat this.io = io; this.encryption = encryption; this.context = context; - this.rowDataReader = new RowDataIteratorReader(tableSchema, + this.rowDataReader = new RowDataFileReader(tableSchema, context.project(), context.nameMapping(), context.caseSensitive()); } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileReader.java similarity index 96% rename from flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java rename to flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileReader.java index 36668bde6b14..4ce22a293970 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataIteratorReader.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileReader.java @@ -45,15 +45,15 @@ import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.util.PartitionUtil; -public class RowDataIteratorReader implements IteratorReader { +public class RowDataFileReader implements FileReader { private final Schema tableSchema; private final Schema projectedSchema; private final String nameMapping; private final boolean caseSensitive; - public RowDataIteratorReader(Schema tableSchema, Schema projectedSchema, - String nameMapping, boolean caseSensitive) { + public RowDataFileReader(Schema tableSchema, Schema projectedSchema, + String nameMapping, boolean caseSensitive) { this.tableSchema = tableSchema; this.projectedSchema = projectedSchema; this.nameMapping = nameMapping; diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java index d21e4a5d1517..401d19d9440b 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java @@ -99,7 +99,7 @@ public static class RewriteMap extends RichMapFunction taskWriterFactory; - private final RowDataIteratorReader rowDataReader; + private final RowDataFileReader rowDataReader; public RewriteMap(Schema schema, String nameMapping, FileIO io, boolean caseSensitive, EncryptionManager encryptionManager, TaskWriterFactory taskWriterFactory) { @@ -109,7 +109,7 @@ public RewriteMap(Schema schema, String nameMapping, FileIO io, boolean caseSens this.caseSensitive = caseSensitive; this.encryptionManager = encryptionManager; this.taskWriterFactory = taskWriterFactory; - this.rowDataReader = new RowDataIteratorReader(schema, schema, nameMapping, caseSensitive); + this.rowDataReader = new RowDataFileReader(schema, schema, nameMapping, caseSensitive); } @Override From 00f18d856203b9d59f06692d6f92579cdf926b67 Mon Sep 17 00:00:00 2001 From: Steven Wu Date: Fri, 13 Aug 2021 13:47:51 -0700 Subject: [PATCH 4/5] fix style error --- .../encryption/InputFilesDecryptor.java | 23 +++++++++++-------- 1 file changed, 14 insertions(+), 9 deletions(-) diff --git a/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java b/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java index 54c48c23764d..f56a9939aaf8 100644 --- a/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java +++ b/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java @@ -1,15 +1,20 @@ /* - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. */ package org.apache.iceberg.encryption; From 7db0ec99281f9810755fd779547adb71f5f0b127 Mon Sep 17 00:00:00 2001 From: Steven Wu Date: Thu, 2 Sep 2021 09:15:57 -0700 Subject: [PATCH 5/5] address openInx's review comments --- .../iceberg/encryption/InputFilesDecryptor.java | 2 +- .../iceberg/flink/source/DataIterator.java | 16 +++++++++------- .../{FileReader.java => FileScanTaskReader.java} | 8 +++++--- .../iceberg/flink/source/FlinkInputFormat.java | 4 ++-- ...eader.java => RowDataFileScanTaskReader.java} | 8 +++++--- .../iceberg/flink/source/RowDataRewriter.java | 4 ++-- 6 files changed, 24 insertions(+), 18 deletions(-) rename flink/src/main/java/org/apache/iceberg/flink/source/{FileReader.java => FileScanTaskReader.java} (86%) rename flink/src/main/java/org/apache/iceberg/flink/source/{RowDataFileReader.java => RowDataFileScanTaskReader.java} (95%) diff --git a/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java b/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java index f56a9939aaf8..6c1e0eb8b250 100644 --- a/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java +++ b/core/src/main/java/org/apache/iceberg/encryption/InputFilesDecryptor.java @@ -45,7 +45,7 @@ public InputFilesDecryptor(CombinedScanTask combinedTask, FileIO io, EncryptionM // decrypt with the batch call to avoid multiple RPCs to a key server, if possible Iterable decryptedFiles = encryption.decrypt(encrypted::iterator); - Map files = Maps.newHashMapWithExpectedSize(combinedTask.files().size()); + Map files = Maps.newHashMapWithExpectedSize(keyMetadata.size()); decryptedFiles.forEach(decrypted -> files.putIfAbsent(decrypted.location(), decrypted)); this.decryptedInputFiles = Collections.unmodifiableMap(files); } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java index 79a44ee884cc..d470b0752304 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/DataIterator.java @@ -22,6 +22,7 @@ import java.io.IOException; import java.io.UncheckedIOException; import java.util.Iterator; +import org.apache.flink.annotation.Internal; import org.apache.iceberg.CombinedScanTask; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.encryption.EncryptionManager; @@ -30,21 +31,22 @@ import org.apache.iceberg.io.FileIO; /** - * Base class of Flink iterators. + * Flink data iterator that reads {@link CombinedScanTask} into a {@link CloseableIterator} * - * @param is the Java class returned by this iterator whose objects contain one or more rows. + * @param is the output data type returned by this iterator. */ +@Internal public class DataIterator implements CloseableIterator { - private final FileReader fileReader; + private final FileScanTaskReader fileScanTaskReader; private final InputFilesDecryptor inputFilesDecryptor; private Iterator tasks; private CloseableIterator currentIterator; - public DataIterator(FileReader fileReader, CombinedScanTask task, + public DataIterator(FileScanTaskReader fileScanTaskReader, CombinedScanTask task, FileIO io, EncryptionManager encryption) { - this.fileReader = fileReader; + this.fileScanTaskReader = fileScanTaskReader; this.inputFilesDecryptor = new InputFilesDecryptor(task, io, encryption); this.tasks = task.files().iterator(); @@ -78,8 +80,8 @@ private void updateCurrentIterator() { } } - private CloseableIterator openTaskIterator(FileScanTask scanTask) throws IOException { - return fileReader.open(scanTask, inputFilesDecryptor); + private CloseableIterator openTaskIterator(FileScanTask scanTask) { + return fileScanTaskReader.open(scanTask, inputFilesDecryptor); } @Override diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/FileReader.java b/flink/src/main/java/org/apache/iceberg/flink/source/FileScanTaskReader.java similarity index 86% rename from flink/src/main/java/org/apache/iceberg/flink/source/FileReader.java rename to flink/src/main/java/org/apache/iceberg/flink/source/FileScanTaskReader.java index 2e339cb4067f..04273016ee2d 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/FileReader.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/FileScanTaskReader.java @@ -20,15 +20,17 @@ package org.apache.iceberg.flink.source; import java.io.Serializable; +import org.apache.flink.annotation.Internal; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.encryption.InputFilesDecryptor; import org.apache.iceberg.io.CloseableIterator; /** * Read a {@link FileScanTask} into a {@link CloseableIterator} + * + * @param is the output data type returned by this iterator. */ -public interface FileReader extends Serializable { - +@Internal +public interface FileScanTaskReader extends Serializable { CloseableIterator open(FileScanTask fileScanTask, InputFilesDecryptor inputFilesDecryptor); - } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java b/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java index d614c3ccaf17..8b757ac31606 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/FlinkInputFormat.java @@ -45,7 +45,7 @@ public class FlinkInputFormat extends RichInputFormat private final FileIO io; private final EncryptionManager encryption; private final ScanContext context; - private final RowDataFileReader rowDataReader; + private final RowDataFileScanTaskReader rowDataReader; private transient DataIterator iterator; private transient long currentReadCount = 0L; @@ -56,7 +56,7 @@ public class FlinkInputFormat extends RichInputFormat this.io = io; this.encryption = encryption; this.context = context; - this.rowDataReader = new RowDataFileReader(tableSchema, + this.rowDataReader = new RowDataFileScanTaskReader(tableSchema, context.project(), context.nameMapping(), context.caseSensitive()); } diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileReader.java b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileScanTaskReader.java similarity index 95% rename from flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileReader.java rename to flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileScanTaskReader.java index 4ce22a293970..fbdb7bf3cc02 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileReader.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataFileScanTaskReader.java @@ -20,6 +20,7 @@ package org.apache.iceberg.flink.source; import java.util.Map; +import org.apache.flink.annotation.Internal; import org.apache.flink.table.data.RowData; import org.apache.iceberg.FileScanTask; import org.apache.iceberg.MetadataColumns; @@ -45,15 +46,16 @@ import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.util.PartitionUtil; -public class RowDataFileReader implements FileReader { +@Internal +public class RowDataFileScanTaskReader implements FileScanTaskReader { private final Schema tableSchema; private final Schema projectedSchema; private final String nameMapping; private final boolean caseSensitive; - public RowDataFileReader(Schema tableSchema, Schema projectedSchema, - String nameMapping, boolean caseSensitive) { + public RowDataFileScanTaskReader(Schema tableSchema, Schema projectedSchema, + String nameMapping, boolean caseSensitive) { this.tableSchema = tableSchema; this.projectedSchema = projectedSchema; this.nameMapping = nameMapping; diff --git a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java index 401d19d9440b..752035e4ea3b 100644 --- a/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java +++ b/flink/src/main/java/org/apache/iceberg/flink/source/RowDataRewriter.java @@ -99,7 +99,7 @@ public static class RewriteMap extends RichMapFunction taskWriterFactory; - private final RowDataFileReader rowDataReader; + private final RowDataFileScanTaskReader rowDataReader; public RewriteMap(Schema schema, String nameMapping, FileIO io, boolean caseSensitive, EncryptionManager encryptionManager, TaskWriterFactory taskWriterFactory) { @@ -109,7 +109,7 @@ public RewriteMap(Schema schema, String nameMapping, FileIO io, boolean caseSens this.caseSensitive = caseSensitive; this.encryptionManager = encryptionManager; this.taskWriterFactory = taskWriterFactory; - this.rowDataReader = new RowDataFileReader(schema, schema, nameMapping, caseSensitive); + this.rowDataReader = new RowDataFileScanTaskReader(schema, schema, nameMapping, caseSensitive); } @Override