This is an automated email from the ASF dual-hosted git repository. zhangyue19921010 pushed a commit to branch stream-binary-copy in repository https://gitbox.apache.org/repos/asf/hudi.git
commit 8ed75bc9882f3f3e2d601bd9d5176a9cc0e1ae3c Author: zhangyue143 <[email protected]> AuthorDate: Tue May 27 17:31:41 2025 +0800 hoodie rewrite handler and clustering using rewrite handler --- .../org/apache/hudi/io/HoodieCreateHandle.java | 14 +- .../apache/hudi/io/HoodieCreateRewriteHandle.java | 114 ++++ .../hudi/io/HoodieUnboundedCreateHandle.java | 2 +- .../hudi/io/SingleFileHandleRewriteFactory.java | 62 ++ ...SparkStreamCopyClusteringExecutionStrategy.java | 148 +++++ .../TestHoodieClientOnCopyOnWriteStorage.java | 8 + .../org/apache/hudi/common/bloom/BloomFilter.java | 6 + .../bloom/HoodieDynamicBoundedBloomFilter.java | 38 ++ .../common/bloom/InternalDynamicBloomFilter.java | 6 +- .../hudi/common/bloom/SimpleBloomFilter.java | 10 + .../apache/hudi/common/config/TypedProperties.java | 2 +- .../storage/rewrite/HoodieFileMetadataMerger.java | 84 +++ .../io/storage/rewrite/HoodieFileRewriter.java | 28 + .../storage/rewrite/HoodieFileRewriterFactory.java | 72 +++ .../io/storage/TestHoodieFileMetadataMerger.java | 195 +++++++ .../hudi/parquet/io/HoodieParquetFileRewriter.java | 131 +++++ .../hudi/parquet/io/HoodieParquetRewriterBase.java | 646 +++++++++++++++++++++ .../parquet/io/HoodieParquetRewriterFactory.java | 54 ++ .../java/org/apache/hudi/parquet/io/TestFile.java | 39 ++ .../apache/hudi/parquet/io/TestFileBuilder.java | 188 ++++++ .../parquet/io/TestHoodieParquetFileRewriter.java | 475 +++++++++++++++ 21 files changed, 2312 insertions(+), 10 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCreateHandle.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCreateHandle.java index 56cb68095a8..e21f18326d2 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCreateHandle.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCreateHandle.java @@ -67,25 +67,25 @@ public class HoodieCreateHandle<T, I, K, O> extends HoodieWriteHandle<T, I, K, O public HoodieCreateHandle(HoodieWriteConfig config, String instantTime, HoodieTable<T, I, K, O> hoodieTable, String partitionPath, String fileId, TaskContextSupplier taskContextSupplier) { this(config, instantTime, hoodieTable, partitionPath, fileId, Option.empty(), - taskContextSupplier, false); + taskContextSupplier, false, true); } public HoodieCreateHandle(HoodieWriteConfig config, String instantTime, HoodieTable<T, I, K, O> hoodieTable, String partitionPath, String fileId, TaskContextSupplier taskContextSupplier, boolean preserveMetadata) { this(config, instantTime, hoodieTable, partitionPath, fileId, Option.empty(), - taskContextSupplier, preserveMetadata); + taskContextSupplier, preserveMetadata, true); } public HoodieCreateHandle(HoodieWriteConfig config, String instantTime, HoodieTable<T, I, K, O> hoodieTable, String partitionPath, String fileId, Option<Schema> overriddenSchema, TaskContextSupplier taskContextSupplier) { - this(config, instantTime, hoodieTable, partitionPath, fileId, overriddenSchema, taskContextSupplier, false); + this(config, instantTime, hoodieTable, partitionPath, fileId, overriddenSchema, taskContextSupplier, false, true); } public HoodieCreateHandle(HoodieWriteConfig config, String instantTime, HoodieTable<T, I, K, O> hoodieTable, String partitionPath, String fileId, Option<Schema> overriddenSchema, - TaskContextSupplier taskContextSupplier, boolean preserveMetadata) { + TaskContextSupplier taskContextSupplier, boolean preserveMetadata, boolean initWriter) { super(config, instantTime, partitionPath, fileId, hoodieTable, overriddenSchema, taskContextSupplier); this.preserveMetadata = preserveMetadata; @@ -103,9 +103,9 @@ public class HoodieCreateHandle<T, I, K, O> extends HoodieWriteHandle<T, I, K, O partitionMetadata.trySave(); createMarkerFile(partitionPath, FSUtils.makeBaseFileName(this.instantTime, this.writeToken, this.fileId, hoodieTable.getBaseFileExtension())); - this.fileWriter = - HoodieFileWriterFactory.getFileWriter(instantTime, path, hoodieTable.getStorage(), config, - writeSchemaWithMetaFields, this.taskContextSupplier, config.getRecordMerger().getRecordType()); + this.fileWriter = initWriter + ? HoodieFileWriterFactory.getFileWriter(instantTime, path, hoodieTable.getStorage(), config, + writeSchemaWithMetaFields, this.taskContextSupplier, config.getRecordMerger().getRecordType()) : null; } catch (IOException e) { throw new HoodieInsertException("Failed to initialize HoodieStorageWriter for path " + path, e); } diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCreateRewriteHandle.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCreateRewriteHandle.java new file mode 100644 index 00000000000..2d54cff6c0f --- /dev/null +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCreateRewriteHandle.java @@ -0,0 +1,114 @@ +/* + * 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.hudi.io; + +import org.apache.hudi.client.WriteStatus; +import org.apache.hudi.common.engine.TaskContextSupplier; +import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.common.util.Option; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.exception.HoodieException; +import org.apache.hudi.exception.HoodieIOException; +import org.apache.hudi.io.storage.rewrite.HoodieFileMetadataMerger; +import org.apache.hudi.io.storage.rewrite.HoodieFileRewriter; +import org.apache.hudi.io.storage.rewrite.HoodieFileRewriterFactory; +import org.apache.hudi.storage.StoragePath; +import org.apache.hudi.table.HoodieTable; + +import org.apache.hadoop.conf.Configuration; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.util.List; + +public class HoodieCreateRewriteHandle<T, I, K, O> extends HoodieCreateHandle<T, I, K, O> { + + private static final Logger LOG = LoggerFactory.getLogger(HoodieCreateRewriteHandle.class); + + protected final HoodieFileRewriter rewriter; + + private final List<StoragePath> inputFiles; + + public HoodieCreateRewriteHandle( + HoodieWriteConfig config, + String instantTime, + String partitionPath, + String fileId, + HoodieTable<T, I, K, O> hoodieTable, + TaskContextSupplier taskContextSupplier, + List<StoragePath> inputFilePaths, + boolean preserveMetadata) { + super( + config, + instantTime, + hoodieTable, + partitionPath, + fileId, + Option.empty(), + taskContextSupplier, + preserveMetadata, + false); + try { + this.inputFiles = inputFilePaths; + HoodieFileMetadataMerger fileMetadataMerger = new HoodieFileMetadataMerger(); + this.rewriter = HoodieFileRewriterFactory.getFileRewriter( + inputFilePaths, + path, + hoodieTable.getStorageConf().unwrapAs(Configuration.class), + hoodieTable.getConfig(), + fileMetadataMerger, + config.getRecordMerger().getRecordType(), this.writeSchemaWithMetaFields); + } catch (IOException e) { + LOG.error("Fail to create file rewriter, cause: ", e); + throw new HoodieException(e); + } + } + + public void rewrite() { + LOG.info("Start to rewrite source files " + this.inputFiles + " into target file: " + this.path); + long start = System.currentTimeMillis(); + long records = 0; + try { + records = this.rewriter.rewrite(); + } catch (IOException e) { + throw new HoodieIOException(e.getMessage(), e); + } finally { + this.recordsWritten = records; + this.insertRecordsWritten = records; + } + LOG.info("Finish rewriting " + this.path + ". Using " + (System.currentTimeMillis() - start) + " mills"); + } + + @Override + public List<WriteStatus> close() { + try { + this.rewriter.close(); + return super.close(); + } catch (IOException e) { + LOG.error("Fail to close the rewrite handle for path: " + path); + throw new HoodieIOException(e.getMessage(), e); + } + } + + @Override + public boolean canWrite(HoodieRecord record) { + return true; + } +} diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieUnboundedCreateHandle.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieUnboundedCreateHandle.java index 19546621e60..4cd55a1f803 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieUnboundedCreateHandle.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieUnboundedCreateHandle.java @@ -43,7 +43,7 @@ public class HoodieUnboundedCreateHandle<T, I, K, O> extends HoodieCreateHandle< String partitionPath, String fileId, TaskContextSupplier taskContextSupplier, boolean preserveHoodieMetadata) { super(config, instantTime, hoodieTable, partitionPath, fileId, Option.empty(), - taskContextSupplier, preserveHoodieMetadata); + taskContextSupplier, preserveHoodieMetadata, true); } @Override diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/SingleFileHandleRewriteFactory.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/SingleFileHandleRewriteFactory.java new file mode 100644 index 00000000000..3d87b66cd73 --- /dev/null +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/SingleFileHandleRewriteFactory.java @@ -0,0 +1,62 @@ +/* + * 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.hudi.io; + +import org.apache.hudi.common.engine.TaskContextSupplier; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.storage.StoragePath; +import org.apache.hudi.table.HoodieTable; + +import java.io.Serializable; +import java.util.List; + +/** + * Rewrite all inputFilePaths related files into one huge file. + */ +public class SingleFileHandleRewriteFactory<T, I, K, O> + extends WriteHandleFactory<T, I, K, O> implements Serializable { + + private final boolean preserveMetadata; + + private final List<StoragePath> inputFilePaths; + + public SingleFileHandleRewriteFactory(List<StoragePath> inputFilePaths, boolean preserveMetadata) { + this.inputFilePaths = inputFilePaths; + this.preserveMetadata = preserveMetadata; + } + + @Override + public HoodieCreateRewriteHandle<T, I, K, O> create( + HoodieWriteConfig config, + String commitTime, + HoodieTable<T, I, K, O> hoodieTable, + String partitionPath, + String fileIdPrefix, + TaskContextSupplier taskContextSupplier) { + return new HoodieCreateRewriteHandle( + config, + commitTime, + partitionPath, + getNextFileId(fileIdPrefix), + hoodieTable, + taskContextSupplier, + inputFilePaths, + preserveMetadata); + } +} \ No newline at end of file diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/SparkStreamCopyClusteringExecutionStrategy.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/SparkStreamCopyClusteringExecutionStrategy.java new file mode 100644 index 00000000000..c2a66ad7ea7 --- /dev/null +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/SparkStreamCopyClusteringExecutionStrategy.java @@ -0,0 +1,148 @@ +/* + * 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.hudi.client.clustering.run.strategy; + +import org.apache.hudi.avro.model.HoodieClusteringPlan; +import org.apache.hudi.client.WriteStatus; +import org.apache.hudi.client.common.HoodieSparkEngineContext; +import org.apache.hudi.common.config.SerializableSchema; +import org.apache.hudi.common.data.HoodieData; +import org.apache.hudi.common.engine.HoodieEngineContext; +import org.apache.hudi.common.engine.TaskContextSupplier; +import org.apache.hudi.common.fs.FSUtils; +import org.apache.hudi.common.model.ClusteringGroupInfo; +import org.apache.hudi.common.model.HoodieFileGroupId; +import org.apache.hudi.common.util.Option; +import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.data.HoodieJavaRDD; +import org.apache.hudi.io.HoodieCreateRewriteHandle; +import org.apache.hudi.io.SingleFileHandleRewriteFactory; +import org.apache.hudi.storage.StoragePath; +import org.apache.hudi.table.HoodieTable; +import org.apache.hudi.table.action.HoodieWriteMetadata; + +import org.apache.avro.Schema; +import org.apache.spark.api.java.JavaRDD; +import org.apache.spark.api.java.JavaSparkContext; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; +import java.util.stream.Stream; +import java.util.stream.StreamSupport; + +import static org.apache.hudi.common.model.HoodieTableType.COPY_ON_WRITE; + +/** + * Clustering strategy to submit single spark jobs using streaming copy + * PAY ATTENTION!!! + * IN THIS STRATEGY + * 1. Only support clustering for cow table. + * 2. Sort function is not supported yet. + * 3. Each clustering group only has one task to write. + */ +public class SparkStreamCopyClusteringExecutionStrategy<T> extends SparkSortAndSizeExecutionStrategy<T> { + + private static final Logger LOG = LoggerFactory.getLogger(SparkStreamCopyClusteringExecutionStrategy.class); + + public SparkStreamCopyClusteringExecutionStrategy( + HoodieTable table, + HoodieEngineContext engineContext, + HoodieWriteConfig writeConfig) { + super(table, engineContext, writeConfig); + } + + @Override + public HoodieWriteMetadata<HoodieData<WriteStatus>> performClustering( + HoodieClusteringPlan clusteringPlan, + Schema schema, + String instantTime) { + if (!supportBinaryStreamCopy()) { + LOG.info("Required conditions for binary stream copy are currently not satisfied, falling back to default clustering behavior"); + return super.performClustering(clusteringPlan, schema, instantTime); + } + LOG.info("Required conditions are currently satisfied, enabling the optimization of using binary stream copy "); + JavaSparkContext engineContext = HoodieSparkEngineContext.getSparkContext(getEngineContext()); + TaskContextSupplier taskContextSupplier = getEngineContext().getTaskContextSupplier(); + SerializableSchema serializableSchema = new SerializableSchema(schema); + boolean shouldPreserveMetadata = Option.ofNullable(clusteringPlan.getPreserveHoodieMetadata()).orElse(false); + List<ClusteringGroupInfo> clusteringGroupInfos = clusteringPlan.getInputGroups() + .stream() + .map(ClusteringGroupInfo::create) + .collect(Collectors.toList()); + + JavaRDD<ClusteringGroupInfo> groupInfoJavaRDD = engineContext.parallelize(clusteringGroupInfos, clusteringGroupInfos.size()); + LOG.info("number of partitions for clustering " + groupInfoJavaRDD.getNumPartitions()); + JavaRDD<WriteStatus> writeStatusRDD = groupInfoJavaRDD + .mapPartitions(clusteringOps -> { + Iterable<ClusteringGroupInfo> clusteringOpsIterable = () -> clusteringOps; + return StreamSupport.stream(clusteringOpsIterable.spliterator(), false) + .flatMap(clusteringOp -> + runClusteringForGroup( + clusteringOp, + clusteringPlan.getStrategy().getStrategyParams(), + shouldPreserveMetadata, + serializableSchema, + taskContextSupplier, + instantTime)) + .iterator(); + }); + + HoodieWriteMetadata<HoodieData<WriteStatus>> writeMetadata = new HoodieWriteMetadata<>(); + writeMetadata.setWriteStatuses(HoodieJavaRDD.of(writeStatusRDD)); + return writeMetadata; + } + + /** + * Submit job to execute clustering for the group. + */ + private Stream<WriteStatus> runClusteringForGroup(ClusteringGroupInfo clusteringOps, Map<String, String> strategyParams, + boolean preserveHoodieMetadata, SerializableSchema schema, + TaskContextSupplier taskContextSupplier, String instantTime) { + List<WriteStatus> statuses = new ArrayList<>(); + List<HoodieFileGroupId> inputFileIds = clusteringOps.getOperations() + .stream() + .map(op -> new HoodieFileGroupId(op.getPartitionPath(), op.getFileId())) + .collect(Collectors.toList()); + List<StoragePath> inputFilePaths = clusteringOps.getOperations() + .stream() + .map(op -> new StoragePath(op.getDataFilePath())) + .collect(Collectors.toList()); + + SingleFileHandleRewriteFactory rewriteFactory = new SingleFileHandleRewriteFactory(inputFilePaths, preserveHoodieMetadata); + HoodieCreateRewriteHandle hoodieRewriteHandle = rewriteFactory.create( + getWriteConfig(), + instantTime, + getHoodieTable(), + inputFileIds.get(0).getPartitionPath(), + FSUtils.createNewFileIdPfx(), + taskContextSupplier); + + hoodieRewriteHandle.rewrite(); + statuses.addAll(hoodieRewriteHandle.close()); + return statuses.stream(); + } + + private boolean supportBinaryStreamCopy() { + return this.getHoodieTable().getMetaClient().getTableType() == COPY_ON_WRITE; + } +} diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java index 4d2bc10a86a..5038338e2b2 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java @@ -1016,6 +1016,14 @@ public class TestHoodieClientOnCopyOnWriteStorage extends HoodieClientTestBase { false, SqlQueryEqualityPreCommitValidator.class.getName(), COUNT_SQL_QUERY_FOR_VALIDATION, ""); } + @Test + public void testStreamingCopyClustering() throws Exception { + String strategy = "org.apache.hudi.client.clustering.run.strategy.SparkStreamCopyClusteringExecutionStrategy"; + HoodieClusteringConfig config = createClusteringBuilder(true, 1).withClusteringExecutionStrategyClass(strategy).build(); + testInsertAndClustering(config, true, true, + false, SqlQueryEqualityPreCommitValidator.class.getName(), COUNT_SQL_QUERY_FOR_VALIDATION, ""); + } + @Test public void testAndValidateClusteringOutputFiles() throws IOException { testAndValidateClusteringOutputFiles(createBrokenClusteringClient(new HoodieException(CLUSTERING_FAILURE)), createClusteringBuilder(true, 2).build(), list2Rdd, rdd2List); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/bloom/BloomFilter.java b/hudi-common/src/main/java/org/apache/hudi/common/bloom/BloomFilter.java index fbc46827dee..7920f30949a 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/bloom/BloomFilter.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/bloom/BloomFilter.java @@ -54,4 +54,10 @@ public interface BloomFilter { * @return the bloom index type code. **/ BloomFilterTypeCode getBloomFilterTypeCode(); + + /** + * Performs a logical OR operations with other BloomFilter. + * @param other + */ + void or(BloomFilter other); } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/bloom/HoodieDynamicBoundedBloomFilter.java b/hudi-common/src/main/java/org/apache/hudi/common/bloom/HoodieDynamicBoundedBloomFilter.java index 5a4381d2ab8..2597fd20f4a 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/bloom/HoodieDynamicBoundedBloomFilter.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/bloom/HoodieDynamicBoundedBloomFilter.java @@ -19,6 +19,7 @@ package org.apache.hudi.common.bloom; import org.apache.hudi.common.util.Base64CodecUtil; +import org.apache.hudi.common.util.ValidationUtils; import org.apache.hudi.exception.HoodieIndexException; import java.io.ByteArrayInputStream; @@ -125,5 +126,42 @@ public class HoodieDynamicBoundedBloomFilter implements BloomFilter { internalDynamicBloomFilter = new InternalDynamicBloomFilter(); internalDynamicBloomFilter.readFields(dis); } + + @Override + public void or(BloomFilter other) { + if (other != null) { + ValidationUtils.checkArgument(other instanceof HoodieDynamicBoundedBloomFilter, + "HoodieDynamicBoundedBloomFilter can only perform OR operations with other HoodieDynamicBoundedBloomFilter."); + HoodieDynamicBoundedBloomFilter otherFilter = (HoodieDynamicBoundedBloomFilter) other; + int sourceMatrixLength = this.getMatrixLength(); + int targetMatrixLength = otherFilter.getMatrixLength(); + if (targetMatrixLength > sourceMatrixLength) { + this.rescaleFromTarget(targetMatrixLength); + } else { + otherFilter.rescaleFromTarget(sourceMatrixLength); + } + this.internalDynamicBloomFilter.or(otherFilter.internalDynamicBloomFilter); + } + } + + /** + * rescale the internal dynamic bloom filter by length + * @param targetMatrixLength + * @return + */ + private HoodieDynamicBoundedBloomFilter rescaleFromTarget(int targetMatrixLength) { + int initMatrixLength = this.internalDynamicBloomFilter.getMatrixLength(); + int needAddRowNum = targetMatrixLength - initMatrixLength; + if (needAddRowNum > 0) { + for (int i = 0; i < needAddRowNum; i++) { + this.internalDynamicBloomFilter.addRow(); + } + } + return this; + } + + private int getMatrixLength() { + return this.internalDynamicBloomFilter.getMatrixLength(); + } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/bloom/InternalDynamicBloomFilter.java b/hudi-common/src/main/java/org/apache/hudi/common/bloom/InternalDynamicBloomFilter.java index 3e068294a0b..cb707238f1a 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/bloom/InternalDynamicBloomFilter.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/bloom/InternalDynamicBloomFilter.java @@ -210,7 +210,7 @@ class InternalDynamicBloomFilter extends InternalFilter { /** * Adds a new row to <i>this</i> dynamic Bloom filter. */ - private void addRow() { + protected void addRow() { InternalBloomFilter[] tmp = new InternalBloomFilter[matrix.length + 1]; System.arraycopy(matrix, 0, tmp, 0, matrix.length); tmp[tmp.length - 1] = new InternalBloomFilter(vectorSize, nbHash, hashType); @@ -236,4 +236,8 @@ class InternalDynamicBloomFilter extends InternalFilter { } return matrix[matrix.length - 1]; } + + public int getMatrixLength() { + return this.matrix.length; + } } \ No newline at end of file diff --git a/hudi-common/src/main/java/org/apache/hudi/common/bloom/SimpleBloomFilter.java b/hudi-common/src/main/java/org/apache/hudi/common/bloom/SimpleBloomFilter.java index c7ada7a54fc..23281cbdaed 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/bloom/SimpleBloomFilter.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/bloom/SimpleBloomFilter.java @@ -19,6 +19,7 @@ package org.apache.hudi.common.bloom; import org.apache.hudi.common.util.Base64CodecUtil; +import org.apache.hudi.common.util.ValidationUtils; import org.apache.hudi.exception.HoodieIndexException; import java.io.ByteArrayInputStream; @@ -155,4 +156,13 @@ public class SimpleBloomFilter implements BloomFilter { private void extractAndSetInternalBloomFilter(DataInputStream dis) throws IOException { this.filter.readFields(dis); } + + @Override + public void or(BloomFilter otherFilter) { + if (otherFilter != null) { + ValidationUtils.checkArgument(otherFilter instanceof SimpleBloomFilter, "SimpleBloomFilter can only perform OR operations with other SimpleBloomFilters."); + SimpleBloomFilter bloomFilter = (SimpleBloomFilter) otherFilter; + this.filter.or(bloomFilter.filter); + } + } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/config/TypedProperties.java b/hudi-common/src/main/java/org/apache/hudi/common/config/TypedProperties.java index 644cc5118be..df7c98104f2 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/config/TypedProperties.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/config/TypedProperties.java @@ -39,7 +39,7 @@ public class TypedProperties extends Properties implements Serializable { super(null); } - protected TypedProperties(Properties defaults) { + public TypedProperties(Properties defaults) { if (Objects.nonNull(defaults)) { for (Enumeration<?> e = defaults.propertyNames(); e.hasMoreElements(); ) { Object k = e.nextElement(); diff --git a/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileMetadataMerger.java b/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileMetadataMerger.java new file mode 100644 index 00000000000..31a6fcc5846 --- /dev/null +++ b/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileMetadataMerger.java @@ -0,0 +1,84 @@ +/* + * 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.hudi.io.storage.rewrite; + +import org.apache.hudi.avro.HoodieBloomFilterWriteSupport; +import org.apache.hudi.common.bloom.BloomFilter; +import org.apache.hudi.common.bloom.BloomFilterFactory; +import org.apache.hudi.common.util.StringUtils; + +import java.util.HashMap; +import java.util.Map; + +import static org.apache.hudi.avro.HoodieBloomFilterWriteSupport.HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY; + +public class HoodieFileMetadataMerger { + + private final HashMap<String, String> mergedMateData = new HashMap<>(); + + private BloomFilter bloomFilter; + + private String minRecordKey; + + private String maxRecordKey; + + public HoodieFileMetadataMerger() { + + } + + public Map<String, String> mergeMetaData(Map<String, String> metaMap) { + mergedMateData.putAll(metaMap); + + String minRecordKey = metaMap.get(HoodieBloomFilterWriteSupport.HOODIE_MIN_RECORD_KEY_FOOTER); + String maxRecordKey = metaMap.get(HoodieBloomFilterWriteSupport.HOODIE_MAX_RECORD_KEY_FOOTER); + if (this.minRecordKey == null || (minRecordKey != null && this.minRecordKey.compareTo(minRecordKey) > 0)) { + this.minRecordKey = minRecordKey; + } + if (this.maxRecordKey == null || (maxRecordKey != null && this.maxRecordKey.compareTo(maxRecordKey) < 0)) { + this.maxRecordKey = maxRecordKey; + } + + String bloomFilterType = metaMap.get(HoodieBloomFilterWriteSupport.HOODIE_BLOOM_FILTER_TYPE_CODE); + String avroBloom = metaMap.get(HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY); + if (!StringUtils.isNullOrEmpty(avroBloom)) { + BloomFilter targetBloomFilter = BloomFilterFactory.fromString(avroBloom, bloomFilterType); + if (this.bloomFilter == null) { + this.bloomFilter = targetBloomFilter; + } else { + this.bloomFilter.or(targetBloomFilter); + } + } + + if (this.minRecordKey != null) { + mergedMateData.put(HoodieBloomFilterWriteSupport.HOODIE_MIN_RECORD_KEY_FOOTER, this.minRecordKey); + } + if (this.maxRecordKey != null) { + mergedMateData.put(HoodieBloomFilterWriteSupport.HOODIE_MAX_RECORD_KEY_FOOTER, this.maxRecordKey); + } + if (this.bloomFilter != null) { + mergedMateData.put(HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY, this.bloomFilter.serializeToString()); + } + + return mergedMateData; + } + + public Map<String, String> getMergedMetaData() { + return mergedMateData; + } +} diff --git a/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileRewriter.java b/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileRewriter.java new file mode 100644 index 00000000000..403496d05dc --- /dev/null +++ b/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileRewriter.java @@ -0,0 +1,28 @@ +/* + * 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.hudi.io.storage.rewrite; + +import java.io.IOException; + +public interface HoodieFileRewriter { + + long rewrite() throws IOException; + + void close() throws IOException; +} diff --git a/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileRewriterFactory.java b/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileRewriterFactory.java new file mode 100644 index 00000000000..a173b31f44d --- /dev/null +++ b/hudi-common/src/main/java/org/apache/hudi/io/storage/rewrite/HoodieFileRewriterFactory.java @@ -0,0 +1,72 @@ +/* + * 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.hudi.io.storage.rewrite; + +import org.apache.hudi.common.config.HoodieConfig; +import org.apache.hudi.common.fs.FSUtils; +import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType; +import org.apache.hudi.common.util.ReflectionUtils; +import org.apache.hudi.exception.HoodieException; +import org.apache.hudi.storage.StoragePath; + +import org.apache.avro.Schema; +import org.apache.hadoop.conf.Configuration; + +import java.io.IOException; +import java.util.List; + +import static org.apache.hudi.common.model.HoodieFileFormat.PARQUET; + +public class HoodieFileRewriterFactory { + + private static HoodieFileRewriterFactory getWriterFactory(HoodieRecordType recordType, String extension) { + if (PARQUET.getFileExtension().equals(extension)) { + try { + Class<?> clazz = ReflectionUtils.getClass("org.apache.hudi.parquet.io.HoodieParquetRewriterFactory"); + return (HoodieFileRewriterFactory) clazz.newInstance(); + } catch (IllegalAccessException | IllegalArgumentException | InstantiationException e) { + throw new HoodieException("Unable to create hoodie avro parquet file writer factory", e); + } + } + throw new UnsupportedOperationException(extension + " file format not supported yet."); + } + + public static <T, I, K, O> HoodieFileRewriter getFileRewriter( + List<StoragePath> inputFilePaths, + StoragePath targetFilePath, + Configuration conf, + HoodieConfig config, + HoodieFileMetadataMerger metadataMerger, + HoodieRecordType recordType, + Schema writeSchemaWithMetaFields) throws IOException { + String extension = FSUtils.getFileExtension(targetFilePath.getName()); + HoodieFileRewriterFactory factory = getWriterFactory(recordType, extension); + return factory.newFileRewriter(inputFilePaths, targetFilePath, conf, config, metadataMerger, writeSchemaWithMetaFields); + } + + protected <T> HoodieFileRewriter newFileRewriter( + List<StoragePath> inputFilePaths, + StoragePath targetFilePath, + Configuration conf, + HoodieConfig config, + HoodieFileMetadataMerger metadataMerger, + Schema writeSchemaWithMetaFields) throws IOException { + throw new UnsupportedOperationException(); + } +} diff --git a/hudi-common/src/test/java/org/apache/hudi/io/storage/TestHoodieFileMetadataMerger.java b/hudi-common/src/test/java/org/apache/hudi/io/storage/TestHoodieFileMetadataMerger.java new file mode 100644 index 00000000000..a5ae5c679a2 --- /dev/null +++ b/hudi-common/src/test/java/org/apache/hudi/io/storage/TestHoodieFileMetadataMerger.java @@ -0,0 +1,195 @@ +/* + * 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.hudi.io.storage; + +import org.apache.hudi.common.bloom.BloomFilter; +import org.apache.hudi.common.bloom.BloomFilterFactory; +import org.apache.hudi.common.bloom.BloomFilterTypeCode; +import org.apache.hudi.common.util.ValidationUtils; +import org.apache.hudi.io.storage.rewrite.HoodieFileMetadataMerger; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.util.HashMap; +import java.util.Map; + +import static org.apache.hudi.avro.HoodieBloomFilterWriteSupport.HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY; +import static org.apache.hudi.avro.HoodieBloomFilterWriteSupport.HOODIE_BLOOM_FILTER_TYPE_CODE; +import static org.apache.hudi.avro.HoodieBloomFilterWriteSupport.HOODIE_MAX_RECORD_KEY_FOOTER; +import static org.apache.hudi.avro.HoodieBloomFilterWriteSupport.HOODIE_MIN_RECORD_KEY_FOOTER; +import static org.junit.jupiter.api.Assertions.assertEquals; + +public class TestHoodieFileMetadataMerger { + + @Test + public void testMinKey() { + HoodieFileMetadataMerger metaMerge = new HoodieFileMetadataMerger(); + // for empty map + Map<String, String> mergeMap = metaMerge.mergeMetaData(newMap()); + Assertions.assertTrue(mergeMap.keySet().size() == 0); + + // just min key + metaMerge = new HoodieFileMetadataMerger(); + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MIN_RECORD_KEY_FOOTER, "1")); + Assertions.assertTrue(mergeMap.keySet().size() == 1); + assertEquals("1", mergeMap.get(HOODIE_MIN_RECORD_KEY_FOOTER)); + + // then with empty map, min key should not change + mergeMap = metaMerge.mergeMetaData(newMap()); + Assertions.assertTrue(mergeMap.keySet().size() == 1); + assertEquals("1", mergeMap.get(HOODIE_MIN_RECORD_KEY_FOOTER)); + + // then with key bigger than 1, should not change + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MIN_RECORD_KEY_FOOTER, "5")); + Assertions.assertTrue(mergeMap.keySet().size() == 1); + assertEquals("1", mergeMap.get(HOODIE_MIN_RECORD_KEY_FOOTER)); + + // then with key smaller than 1, should change + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MIN_RECORD_KEY_FOOTER, "0")); + Assertions.assertTrue(mergeMap.keySet().size() == 1); + assertEquals("0", mergeMap.get(HOODIE_MIN_RECORD_KEY_FOOTER)); + } + + @Test + public void testMaxKey() { + HoodieFileMetadataMerger metaMerge = new HoodieFileMetadataMerger(); + // for empty map + Map<String, String> mergeMap = metaMerge.mergeMetaData(newMap()); + Assertions.assertTrue(mergeMap.keySet().size() == 0); + + // just max key + metaMerge = new HoodieFileMetadataMerger(); + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MAX_RECORD_KEY_FOOTER, "1")); + Assertions.assertTrue(mergeMap.keySet().size() == 1); + assertEquals("1", mergeMap.get(HOODIE_MAX_RECORD_KEY_FOOTER)); + + // then with empty map, min key should not change + mergeMap = metaMerge.mergeMetaData(newMap()); + Assertions.assertTrue(mergeMap.keySet().size() == 1); + assertEquals("1", mergeMap.get(HOODIE_MAX_RECORD_KEY_FOOTER)); + + // then with key smaller than 1, should not change + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MAX_RECORD_KEY_FOOTER, "0")); + Assertions.assertTrue(mergeMap.keySet().size() == 1); + assertEquals("1", mergeMap.get(HOODIE_MAX_RECORD_KEY_FOOTER)); + + // then with key bigger than 1, should change + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MAX_RECORD_KEY_FOOTER, "5")); + assertEquals("5", mergeMap.get(HOODIE_MAX_RECORD_KEY_FOOTER)); + Assertions.assertFalse(mergeMap.containsKey(HOODIE_MIN_RECORD_KEY_FOOTER)); + } + + @Test + public void tesMaxAndMin() { + HoodieFileMetadataMerger metaMerge = new HoodieFileMetadataMerger(); + + Map<String, String> mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MIN_RECORD_KEY_FOOTER, "1", HOODIE_MAX_RECORD_KEY_FOOTER, "6")); + Assertions.assertTrue(mergeMap.keySet().size() == 2); + assertEquals("1", mergeMap.get(HOODIE_MIN_RECORD_KEY_FOOTER)); + assertEquals("6", mergeMap.get(HOODIE_MAX_RECORD_KEY_FOOTER)); + + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MIN_RECORD_KEY_FOOTER, "0", HOODIE_MAX_RECORD_KEY_FOOTER, "5")); + Assertions.assertTrue(mergeMap.keySet().size() == 2); + assertEquals("0", mergeMap.get(HOODIE_MIN_RECORD_KEY_FOOTER)); + assertEquals("6", mergeMap.get(HOODIE_MAX_RECORD_KEY_FOOTER)); + + mergeMap = metaMerge.mergeMetaData(newMap(HOODIE_MIN_RECORD_KEY_FOOTER, "4", HOODIE_MAX_RECORD_KEY_FOOTER, "5")); + Assertions.assertTrue(mergeMap.keySet().size() == 2); + assertEquals("0", mergeMap.get(HOODIE_MIN_RECORD_KEY_FOOTER)); + assertEquals("6", mergeMap.get(HOODIE_MAX_RECORD_KEY_FOOTER)); + } + + @ParameterizedTest() + @ValueSource(strings = {"SIMPLE", "DYNAMIC_V0"}) + public void testBloomFilter(String bloomFilterType) { + HoodieFileMetadataMerger metaMerge = new HoodieFileMetadataMerger(); + int[] sizes = {100, 1000, 10000}; + BloomFilter bloomFilter = null; + for (int size : sizes) { + BloomFilter filter = getBloomFilter(bloomFilterType, 10000, 0.000001, 100000); + for (int i = 0; i < size; i++) { + String key = String.format("key%d", size + i); + filter.add(key); + } + if (bloomFilter == null) { + bloomFilter = filter; + } else { + bloomFilter.or(filter); + } + + metaMerge.mergeMetaData( + newMap( + HOODIE_BLOOM_FILTER_TYPE_CODE, bloomFilterType, + HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY, filter.serializeToString()) + ); + } + + Map<String, String> mergedMetaData = metaMerge.getMergedMetaData(); + assertEquals(bloomFilterType, mergedMetaData.get(HOODIE_BLOOM_FILTER_TYPE_CODE)); + assertEquals(bloomFilter.serializeToString(), mergedMetaData.get(HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY)); + } + + @Test + public void testDifferentTypeOfBloomFilter() { + HoodieFileMetadataMerger metaMerge = new HoodieFileMetadataMerger(); + BloomFilter simpleFilter = getBloomFilter(BloomFilterTypeCode.SIMPLE.name(), 10000, 0.000001, 100000); + for (int i = 0; i < 100; i++) { + String key = String.format("key%d", 100 + i); + simpleFilter.add(key); + } + metaMerge.mergeMetaData( + newMap( + HOODIE_BLOOM_FILTER_TYPE_CODE, BloomFilterTypeCode.SIMPLE.name(), + HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY, simpleFilter.serializeToString()) + ); + BloomFilter dynamicFilter = getBloomFilter(BloomFilterTypeCode.DYNAMIC_V0.name(), 10000, 0.000001, 100000); + for (int i = 0; i < 100; i++) { + String key = String.format("key%d", 100 + i); + dynamicFilter.add(key); + } + + Assertions.assertThrows(IllegalArgumentException.class, () -> + metaMerge.mergeMetaData( + newMap( + HOODIE_BLOOM_FILTER_TYPE_CODE, BloomFilterTypeCode.DYNAMIC_V0.name(), + HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY, dynamicFilter.serializeToString()) + ) + ); + } + + private BloomFilter getBloomFilter(String typeCode, int numEntries, double errorRate, int maxEntries) { + if (typeCode.equalsIgnoreCase(BloomFilterTypeCode.SIMPLE.name())) { + return BloomFilterFactory.createBloomFilter(numEntries, errorRate, -1, typeCode); + } else { + return BloomFilterFactory.createBloomFilter(numEntries, errorRate, maxEntries, typeCode); + } + } + + private Map<String, String> newMap(String... kvs) { + ValidationUtils.checkArgument(kvs.length == 0 || kvs.length % 2 == 0, "num of input args should be 0 or multiples of 2"); + HashMap map = new HashMap(); + for (int i = 0; i < kvs.length; i += 2) { + map.put(kvs[i], kvs[i + 1]); + } + return map; + } +} \ No newline at end of file diff --git a/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetFileRewriter.java b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetFileRewriter.java new file mode 100644 index 00000000000..849394db507 --- /dev/null +++ b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetFileRewriter.java @@ -0,0 +1,131 @@ +/* + * 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.hudi.parquet.io; + +import org.apache.hudi.io.storage.HoodieParquetConfig; +import org.apache.hudi.io.storage.rewrite.HoodieFileMetadataMerger; +import org.apache.hudi.io.storage.rewrite.HoodieFileRewriter; +import org.apache.hudi.storage.StoragePath; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.parquet.HadoopReadOptions; +import org.apache.parquet.Preconditions; +import org.apache.parquet.hadoop.metadata.BlockMetaData; +import org.apache.parquet.hadoop.metadata.FileMetaData; +import org.apache.parquet.hadoop.util.CompressionConverter; +import org.apache.parquet.hadoop.util.HadoopInputFile; +import org.apache.parquet.schema.MessageType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import java.util.Queue; +import java.util.Set; + +public class HoodieParquetFileRewriter extends HoodieParquetRewriterBase implements HoodieFileRewriter { + + private static final Logger LOG = LoggerFactory.getLogger(HoodieParquetFileRewriter.class); + + // Reader and relevant states of the in-processing input file + private Queue<CompressionConverter.TransParquetFileReader> inputFiles = new LinkedList<>(); + + private Map<String, String> extraMetaData = new HashMap<>(); + + // The reader for the current input file + private CompressionConverter.TransParquetFileReader reader = null; + + private HoodieFileMetadataMerger metadataMerger; + + public HoodieParquetFileRewriter( + List<StoragePath> inputFilesPath, + StoragePath outPutFile, + Configuration conf, + HoodieParquetConfig parquetConfig, + HoodieFileMetadataMerger metadataMerger, + MessageType writeSchema) { + super(new Path(outPutFile.toUri()), conf); + this.metadataMerger = metadataMerger; + openInputFiles(inputFilesPath, conf); + initFileWriter(parquetConfig, writeSchema); + initNextReader(); + } + + @Override + protected Map<String, String> finalizeMetadata() { + return this.extraMetaData; + } + + @Override + public long rewrite() throws IOException { + Set<String> allOriginalCreatedBys = new HashSet<>(); + while (reader != null) { + List<BlockMetaData> rowGroups = reader.getRowGroups(); + FileMetaData fileMetaData = reader.getFooter().getFileMetaData(); + String createdBy = fileMetaData.getCreatedBy(); + allOriginalCreatedBys.add(createdBy); + Map<String, String> metaMap = fileMetaData.getKeyValueMetaData(); + metadataMerger.mergeMetaData(metaMap); + + for (BlockMetaData block : rowGroups) { + processBlocksFromReader(reader, reader.readNextRowGroup(), block, createdBy); + } + initNextReader(); + } + extraMetaData.putAll(metadataMerger.getMergedMetaData()); + extraMetaData.put(ORIGINAL_CREATED_BY_KEY, String.join("\n", allOriginalCreatedBys)); + return totalRecordsWritten; + } + + // Open all input files to validate their schemas are compatible to merge + private void openInputFiles(List<StoragePath> inputFiles, Configuration conf) { + Preconditions.checkArgument(inputFiles != null && !inputFiles.isEmpty(), "No input files"); + + for (StoragePath inputFile : inputFiles) { + try { + CompressionConverter.TransParquetFileReader reader = new CompressionConverter.TransParquetFileReader( + HadoopInputFile.fromPath(new Path(inputFile.toUri()), conf), + HadoopReadOptions.builder(conf).build()); + this.inputFiles.add(reader); + } catch (IOException e) { + throw new IllegalArgumentException("Failed to open input file: " + inputFile, e); + } + } + } + + // Routines to get reader of next input file and set up relevant states + private void initNextReader() { + if (reader != null) { + LOG.info("Finish rewriting input file: {}", reader.getFile()); + } + + if (inputFiles.isEmpty()) { + reader = null; + return; + } + + reader = inputFiles.poll(); + LOG.info("Rewriting input file: {}, remaining files: {}", reader.getFile(), inputFiles.size()); + } +} diff --git a/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetRewriterBase.java b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetRewriterBase.java new file mode 100644 index 00000000000..7418eca2797 --- /dev/null +++ b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetRewriterBase.java @@ -0,0 +1,646 @@ +/* + * 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.hudi.parquet.io; + +import org.apache.hudi.common.model.HoodieRecord; +import org.apache.hudi.exception.HoodieException; +import org.apache.hudi.io.storage.HoodieParquetConfig; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.parquet.bytes.BytesInput; +import org.apache.parquet.column.ColumnDescriptor; +import org.apache.parquet.column.ColumnReader; +import org.apache.parquet.column.ColumnWriteStore; +import org.apache.parquet.column.ColumnWriter; +import org.apache.parquet.column.EncodingStats; +import org.apache.parquet.column.ParquetProperties; +import org.apache.parquet.column.impl.ColumnReadStoreImpl; +import org.apache.parquet.column.page.DictionaryPage; +import org.apache.parquet.column.page.PageReadStore; +import org.apache.parquet.column.statistics.Statistics; +import org.apache.parquet.column.values.bloomfilter.BloomFilter; +import org.apache.parquet.compression.CompressionCodecFactory; +import org.apache.parquet.crypto.FileEncryptionProperties; +import org.apache.parquet.format.BlockCipher; +import org.apache.parquet.format.DataPageHeader; +import org.apache.parquet.format.DataPageHeaderV2; +import org.apache.parquet.format.DictionaryPageHeader; +import org.apache.parquet.format.PageHeader; +import org.apache.parquet.format.converter.ParquetMetadataConverter; +import org.apache.parquet.hadoop.CodecFactory; +import org.apache.parquet.hadoop.ColumnChunkPageWriteStore; +import org.apache.parquet.hadoop.ParquetFileWriter; +import org.apache.parquet.hadoop.metadata.BlockMetaData; +import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData; +import org.apache.parquet.hadoop.metadata.ColumnPath; +import org.apache.parquet.hadoop.metadata.CompressionCodecName; +import org.apache.parquet.hadoop.metadata.ParquetMetadata; +import org.apache.parquet.hadoop.util.CompressionConverter; +import org.apache.parquet.hadoop.util.HadoopCodecs; +import org.apache.parquet.hadoop.util.HadoopOutputFile; +import org.apache.parquet.internal.column.columnindex.ColumnIndex; +import org.apache.parquet.internal.column.columnindex.OffsetIndex; +import org.apache.parquet.io.ParquetEncodingException; +import org.apache.parquet.io.api.Binary; +import org.apache.parquet.io.api.Converter; +import org.apache.parquet.io.api.GroupConverter; +import org.apache.parquet.io.api.PrimitiveConverter; +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.LogicalTypeAnnotation; +import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.Type; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.Closeable; +import java.io.IOException; +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +import static org.apache.parquet.column.ParquetProperties.DEFAULT_COLUMN_INDEX_TRUNCATE_LENGTH; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_STATISTICS_TRUNCATE_LENGTH; +import static org.apache.parquet.hadoop.ParquetWriter.MAX_PADDING_SIZE_DEFAULT; + +/** + * Copy from parquet-hadoop 1.13.1 org.apache.parquet.hadoop.rewrite.ParquetRewriter + * The reason we did not extends ParquetRewriter is + * 1. we need to control copy operation at block level + * 2. We need to handle schema evolution + * 3. We need to combine file metas added by hudi, such as 'hoodie_min_record_key'/'hoodie_max_record_key'/'org.apache.hudi.bloomfilter' + * 4. We need to overwrite column '_hoodie_file_name' with the output file name + */ +public abstract class HoodieParquetRewriterBase implements Closeable { + + private static final Logger LOG = LoggerFactory.getLogger(HoodieParquetRewriterBase.class); + + // Key to store original writer version in the file key-value metadata + public static final String ORIGINAL_CREATED_BY_KEY = "original.created.by"; + + private final int pageBufferSize = ParquetProperties.DEFAULT_PAGE_SIZE * 2; + + private final byte[] pageBuffer = new byte[pageBufferSize]; + + private CompressionCodecName newCodecName = null; + + private Map<ColumnPath, Binary> maskColumns = new HashMap<>(); + + // Writer to rewrite the input files + private ParquetFileWriter writer; + + // Number of blocks written which is used to keep track of the actual row group ordinal + private int numBlocksRewritten = 0; + + protected long totalRecordsWritten = 0; + + // Schema of input files (should be the same) and to write to the output file + protected MessageType requiredSchema = null; + + private Path outPutFile; + + protected Configuration conf; + + public HoodieParquetRewriterBase(Path outPutFile, Configuration conf) { + this.outPutFile = outPutFile; + this.conf = conf; + + // For meta column '_hoodie_file_name', rewriter will mask value with output file name + Binary maskValue = Binary.fromString(outPutFile.getName()); + maskColumns.put(ColumnPath.fromDotString(HoodieRecord.FILENAME_METADATA_FIELD), maskValue); + } + + protected void initFileWriter(HoodieParquetConfig config, MessageType schema) { + try { + this.requiredSchema = schema; + this.newCodecName = config.getCompressionCodecName(); + ParquetFileWriter.Mode writerMode = ParquetFileWriter.Mode.CREATE; + writer = new ParquetFileWriter( + HadoopOutputFile.fromPath(outPutFile, conf), + schema, + writerMode, + config.getBlockSize(), + MAX_PADDING_SIZE_DEFAULT, + DEFAULT_COLUMN_INDEX_TRUNCATE_LENGTH, + DEFAULT_STATISTICS_TRUNCATE_LENGTH, + ParquetProperties.DEFAULT_PAGE_WRITE_CHECKSUM_ENABLED, + (FileEncryptionProperties) null); + writer.start(); + LOG.info("init rewriter "); + } catch (Exception e) { + LOG.error("failed to init parquetRewriter", e); + throw new HoodieException(e); + } + } + + @Override + public void close() throws IOException { + Map<String, String> extraMetaData = finalizeMetadata(); + writer.end(extraMetaData == null ? new HashMap<>() : extraMetaData); + } + + protected abstract Map<String, String> finalizeMetadata(); + + public void processBlocksFromReader( + CompressionConverter.TransParquetFileReader reader, + PageReadStore store, + BlockMetaData block, + String originalCreatedBy) throws IOException { + if (store == null) { + LOG.info("stores is empty"); + return; + } + + totalRecordsWritten += store.getRowCount(); + ColumnReadStoreImpl crStore = new ColumnReadStoreImpl( + store, + new DummyGroupConverter(), + requiredSchema, + originalCreatedBy); + Map<ColumnPath, ColumnDescriptor> descriptorsMap = requiredSchema.getColumns() + .stream() + .collect(Collectors.toMap(x -> ColumnPath.get(x.getPath()), x -> x)); + + writer.startBlock(store.getRowCount()); + List<ColumnChunkMetaData> columnsInOrder = block.getColumns(); + + for (int i = 0, columnId = 0; i < columnsInOrder.size(); i++) { + ColumnChunkMetaData chunk = columnsInOrder.get(i); + ColumnDescriptor descriptor = descriptorsMap.get(chunk.getPath()); + + // This column has been pruned. + if (descriptor == null) { + continue; + } + + // If a column is encrypted, we simply throw exception. + // Later we can add a feature to trans-encrypt it with different keys + if (chunk.isEncrypted()) { + throw new IOException("Column " + chunk.getPath().toDotString() + " is already encrypted"); + } + reader.setStreamPosition(chunk.getStartingPos()); + CompressionCodecName newCodecName = this.newCodecName == null ? chunk.getCodec() : this.newCodecName; + + if (maskColumns != null && maskColumns.containsKey(chunk.getPath())) { + // Mask column and compress it again. + Binary maskValue = maskColumns.get(chunk.getPath()); + if (maskValue != null) { + maskColumn( + descriptor, + chunk, + crStore, + writer, + requiredSchema, + newCodecName, + maskValue); + } + } else if (this.newCodecName != null && this.newCodecName != chunk.getCodec()) { + // Translate compression and/or encryption + writer.startColumn(descriptor, crStore.getColumnReader(descriptor).getTotalValueCount(), newCodecName); + processChunk(reader, chunk, newCodecName, false, originalCreatedBy); + writer.endColumn(); + } else { + // Nothing changed, simply copy the binary data. + BloomFilter bloomFilter = reader.readBloomFilter(chunk); + ColumnIndex columnIndex = reader.readColumnIndex(chunk); + OffsetIndex offsetIndex = reader.readOffsetIndex(chunk); + writer.appendColumnChunk(descriptor, reader.getStream(), chunk, bloomFilter, columnIndex, offsetIndex); + } + + columnId++; + } + + // append missed columns + ParquetMetadata meta = reader.getFooter(); + ColumnChunkMetaData columnChunkMetaData = columnsInOrder.get(0); + EncodingStats encodingStats = columnChunkMetaData.getEncodingStats(); + List<ColumnDescriptor> missedColumns = missedColumns(requiredSchema, meta.getFileMetaData().getSchema()); + for (ColumnDescriptor descriptor : missedColumns) { + addNullColumn( + descriptor, + store.getRowCount(), + encodingStats, + writer, + requiredSchema, + newCodecName); + } + + writer.endBlock(); + numBlocksRewritten++; + } + + private void processChunk( + CompressionConverter.TransParquetFileReader reader, + ColumnChunkMetaData chunk, + CompressionCodecName newCodecName, + boolean encryptColumn, + String originalCreatedBy) throws IOException { + CompressionCodecFactory codecFactory = HadoopCodecs.newFactory(0); + CompressionCodecFactory.BytesInputDecompressor decompressor = null; + CompressionCodecFactory.BytesInputCompressor compressor = null; + if (!newCodecName.equals(chunk.getCodec())) { + // Re-compress only if a different codec has been specified + decompressor = codecFactory.getDecompressor(chunk.getCodec()); + compressor = codecFactory.getCompressor(newCodecName); + } + + // EncryptorRunTime is only provided when encryption is required + BlockCipher.Encryptor metaEncryptor = null; + BlockCipher.Encryptor dataEncryptor = null; + byte[] dictPageAAD = null; + byte[] dataPageAAD = null; + byte[] dictPageHeaderAAD = null; + byte[] dataPageHeaderAAD = null; + + ColumnIndex columnIndex = reader.readColumnIndex(chunk); + OffsetIndex offsetIndex = reader.readOffsetIndex(chunk); + + reader.setStreamPosition(chunk.getStartingPos()); + DictionaryPage dictionaryPage = null; + long readValues = 0; + Statistics statistics = null; + ParquetMetadataConverter converter = new ParquetMetadataConverter(); + int pageOrdinal = 0; + long totalChunkValues = chunk.getValueCount(); + while (readValues < totalChunkValues) { + PageHeader pageHeader = reader.readPageHeader(); + int compressedPageSize = pageHeader.getCompressed_page_size(); + byte[] pageLoad; + switch (pageHeader.type) { + case DICTIONARY_PAGE: + if (dictionaryPage != null) { + throw new IOException("has more than one dictionary page in column chunk"); + } + //No quickUpdatePageAAD needed for dictionary page + DictionaryPageHeader dictPageHeader = pageHeader.dictionary_page_header; + pageLoad = processPageLoad(reader, + true, + compressor, + decompressor, + pageHeader.getCompressed_page_size(), + pageHeader.getUncompressed_page_size(), + encryptColumn, + dataEncryptor, + dictPageAAD); + writer.writeDictionaryPage(new DictionaryPage(BytesInput.from(pageLoad), + pageHeader.getUncompressed_page_size(), + dictPageHeader.getNum_values(), + converter.getEncoding(dictPageHeader.getEncoding())), + metaEncryptor, + dictPageHeaderAAD); + break; + case DATA_PAGE: + DataPageHeader headerV1 = pageHeader.data_page_header; + pageLoad = processPageLoad(reader, + true, + compressor, + decompressor, + pageHeader.getCompressed_page_size(), + pageHeader.getUncompressed_page_size(), + encryptColumn, + dataEncryptor, + dataPageAAD); + statistics = convertStatistics( + originalCreatedBy, chunk.getPrimitiveType(), headerV1.getStatistics(), columnIndex, pageOrdinal, converter); + readValues += headerV1.getNum_values(); + if (offsetIndex != null) { + long rowCount = 1 + offsetIndex.getLastRowIndex( + pageOrdinal, totalChunkValues) - offsetIndex.getFirstRowIndex(pageOrdinal); + writer.writeDataPage(toIntWithCheck(headerV1.getNum_values()), + pageHeader.getUncompressed_page_size(), + BytesInput.from(pageLoad), + statistics, + toIntWithCheck(rowCount), + converter.getEncoding(headerV1.getRepetition_level_encoding()), + converter.getEncoding(headerV1.getDefinition_level_encoding()), + converter.getEncoding(headerV1.getEncoding()), + metaEncryptor, + dataPageHeaderAAD); + } else { + writer.writeDataPage(toIntWithCheck(headerV1.getNum_values()), + pageHeader.getUncompressed_page_size(), + BytesInput.from(pageLoad), + statistics, + converter.getEncoding(headerV1.getRepetition_level_encoding()), + converter.getEncoding(headerV1.getDefinition_level_encoding()), + converter.getEncoding(headerV1.getEncoding()), + metaEncryptor, + dataPageHeaderAAD); + } + pageOrdinal++; + break; + case DATA_PAGE_V2: + DataPageHeaderV2 headerV2 = pageHeader.data_page_header_v2; + int rlLength = headerV2.getRepetition_levels_byte_length(); + BytesInput rlLevels = readBlockAllocate(rlLength, reader); + int dlLength = headerV2.getDefinition_levels_byte_length(); + BytesInput dlLevels = readBlockAllocate(dlLength, reader); + int payLoadLength = pageHeader.getCompressed_page_size() - rlLength - dlLength; + int rawDataLength = pageHeader.getUncompressed_page_size() - rlLength - dlLength; + pageLoad = processPageLoad( + reader, + headerV2.is_compressed, + compressor, + decompressor, + payLoadLength, + rawDataLength, + encryptColumn, + dataEncryptor, + dataPageAAD); + statistics = convertStatistics( + originalCreatedBy, chunk.getPrimitiveType(), headerV2.getStatistics(), columnIndex, pageOrdinal, converter); + readValues += headerV2.getNum_values(); + writer.writeDataPageV2(headerV2.getNum_rows(), + headerV2.getNum_nulls(), + headerV2.getNum_values(), + rlLevels, + dlLevels, + converter.getEncoding(headerV2.getEncoding()), + BytesInput.from(pageLoad), + rawDataLength, + statistics); + pageOrdinal++; + break; + default: + LOG.debug("skipping page of type {} of size {}", pageHeader.getType(), compressedPageSize); + break; + } + } + } + + private Statistics convertStatistics( + String createdBy, + PrimitiveType type, + org.apache.parquet.format.Statistics pageStatistics, + ColumnIndex columnIndex, + int pageIndex, + ParquetMetadataConverter converter) throws IOException { + if (columnIndex != null) { + if (columnIndex.getNullPages() == null) { + throw new IOException("columnIndex has null variable 'nullPages' which indicates corrupted data for type: " + + type.getName()); + } + if (pageIndex > columnIndex.getNullPages().size()) { + throw new IOException("There are more pages " + pageIndex + " found in the column than in the columnIndex " + + columnIndex.getNullPages().size()); + } + Statistics.Builder statsBuilder = + Statistics.getBuilderForReading(type); + statsBuilder.withNumNulls(columnIndex.getNullCounts().get(pageIndex)); + + if (!columnIndex.getNullPages().get(pageIndex)) { + statsBuilder.withMin(columnIndex.getMinValues().get(pageIndex).array().clone()); + statsBuilder.withMax(columnIndex.getMaxValues().get(pageIndex).array().clone()); + } + return statsBuilder.build(); + } else if (pageStatistics != null) { + return converter.fromParquetStatistics(createdBy, pageStatistics, type); + } else { + return null; + } + } + + private byte[] processPageLoad( + CompressionConverter.TransParquetFileReader reader, + boolean isCompressed, + CompressionCodecFactory.BytesInputCompressor compressor, + CompressionCodecFactory.BytesInputDecompressor decompressor, + int payloadLength, + int rawDataLength, + boolean encrypt, + BlockCipher.Encryptor dataEncryptor, + byte[] add) throws IOException { + BytesInput data = readBlock(payloadLength, reader); + + // recompress page load + if (compressor != null) { + if (isCompressed) { + data = decompressor.decompress(data, rawDataLength); + } + data = compressor.compress(data); + } + + if (!encrypt) { + return data.toByteArray(); + } + + // encrypt page load + return dataEncryptor.encrypt(data.toByteArray(), add); + } + + public BytesInput readBlock(int length, CompressionConverter.TransParquetFileReader reader) throws IOException { + byte[] data; + if (length > pageBufferSize) { + data = new byte[length]; + } else { + data = pageBuffer; + } + reader.blockRead(data, 0, length); + return BytesInput.from(data, 0, length); + } + + public BytesInput readBlockAllocate(int length, CompressionConverter.TransParquetFileReader reader) throws IOException { + byte[] data = new byte[length]; + reader.blockRead(data, 0, length); + return BytesInput.from(data, 0, length); + } + + private int toIntWithCheck(long size) { + if ((int) size != size) { + throw new ParquetEncodingException("size is bigger than " + Integer.MAX_VALUE + " bytes: " + size); + } + return (int) size; + } + + private void maskColumn( + ColumnDescriptor descriptor, + ColumnChunkMetaData chunk, + ColumnReadStoreImpl crStore, + ParquetFileWriter writer, + MessageType schema, + CompressionCodecName newCodecName, + Binary maskValue) throws IOException { + + long totalChunkValues = chunk.getValueCount(); + ColumnReader cReader = crStore.getColumnReader(descriptor); + + ParquetProperties.WriterVersion writerVersion = chunk.getEncodingStats().usesV2Pages() + ? ParquetProperties.WriterVersion.PARQUET_2_0 : ParquetProperties.WriterVersion.PARQUET_1_0; + ParquetProperties props = ParquetProperties.builder() + .withWriterVersion(writerVersion) + .build(); + CodecFactory codecFactory = new CodecFactory(new Configuration(), props.getPageSizeThreshold()); + CodecFactory.BytesCompressor compressor = codecFactory.getCompressor(newCodecName); + + // Create new schema that only has the current column + MessageType newSchema = newSchema(schema, descriptor); + ColumnChunkPageWriteStore cPageStore = new ColumnChunkPageWriteStore( + compressor, newSchema, props.getAllocator(), props.getColumnIndexTruncateLength(), + props.getPageWriteChecksumEnabled(), writer.getEncryptor(), numBlocksRewritten); + ColumnWriteStore cStore = props.newColumnWriteStore(newSchema, cPageStore); + ColumnWriter cWriter = cStore.getColumnWriter(descriptor); + + for (int i = 0; i < totalChunkValues; i++) { + int rlvl = cReader.getCurrentRepetitionLevel(); + int dlvl = cReader.getCurrentDefinitionLevel(); + cWriter.write(maskValue, rlvl, dlvl); + cStore.endRecord(); + } + + cStore.flush(); + cPageStore.flushToFileWriter(writer); + + cStore.close(); + cWriter.close(); + } + + private void addNullColumn( + ColumnDescriptor descriptor, + long totalChunkValues, + EncodingStats encodingStats, + ParquetFileWriter writer, + MessageType schema, + CompressionCodecName newCodecName) throws IOException { + + ParquetProperties.WriterVersion writerVersion = encodingStats.usesV2Pages() + ? ParquetProperties.WriterVersion.PARQUET_2_0 : ParquetProperties.WriterVersion.PARQUET_1_0; + ParquetProperties props = ParquetProperties.builder() + .withWriterVersion(writerVersion) + .build(); + CodecFactory codecFactory = new CodecFactory(new Configuration(), props.getPageSizeThreshold()); + CodecFactory.BytesCompressor compressor = codecFactory.getCompressor(newCodecName); + + // Create new schema that only has the current column + MessageType newSchema = newSchema(schema, descriptor); + ColumnChunkPageWriteStore cPageStore = new ColumnChunkPageWriteStore( + compressor, + newSchema, + props.getAllocator(), + props.getColumnIndexTruncateLength(), + props.getPageWriteChecksumEnabled(), + writer.getEncryptor(), + numBlocksRewritten); + ColumnWriteStore cStore = props.newColumnWriteStore(newSchema, cPageStore); + ColumnWriter cWriter = cStore.getColumnWriter(descriptor); + int dMax = descriptor.getMaxDefinitionLevel(); + + for (int i = 0; i < totalChunkValues; i++) { + int rlvl = 0; + int dlvl = 0; + if (dlvl == dMax) { + // since we checked ether optional or repeated, dlvl should be > 0 + if (dlvl == 0) { + throw new IOException("definition level is detected to be 0 for column " + + Arrays.stream(descriptor.getPath()).collect(Collectors.joining(".")) + " to be nullified"); + } + // we just write one null for the whole list at the top level, + // instead of nullify the elements in the list one by one + if (rlvl == 0) { + cWriter.writeNull(rlvl, dlvl - 1); + } + } else { + cWriter.writeNull(rlvl, dlvl); // 因为repeatition level没有重复所以后面都是以0在第一层,definition level是字段path的第0层 + } + cStore.endRecord(); + } + + cStore.flush(); + cPageStore.flushToFileWriter(writer); + + cStore.close(); + cWriter.close(); + } + + private List<ColumnDescriptor> missedColumns(MessageType requiredSchema, MessageType fileSchema) { + return requiredSchema.getColumns().stream() + .filter(col -> !fileSchema.containsPath(col.getPath())) + .collect(Collectors.toList()); + } + + private MessageType newSchema(MessageType schema, ColumnDescriptor descriptor) { + String[] path = descriptor.getPath(); + Type type = schema.getType(path); + if (path.length == 1) { + return new MessageType(schema.getName(), type); + } + + for (Type field : schema.getFields()) { + if (!field.isPrimitive() && path[0].equals(field.getName())) { + Type newType = extractField(field.asGroupType(), type); + if (newType != null) { + if (LogicalTypeAnnotation.mapType().equals(field.getLogicalTypeAnnotation())) { + return new MessageType(schema.getName(), new MessageType(field.getName(), newType)); + } else { + return new MessageType(schema.getName(), newType); + } + } + } + } + + // We should never hit this because 'type' is returned by schema.getType(). + throw new RuntimeException("No field is found"); + } + + private Type extractField(GroupType candidate, Type targetField) { + if (targetField.equals(candidate)) { + return targetField; + } + + // In case 'type' is a descendants of candidate + for (Type field : candidate.asGroupType().getFields()) { + if (field.isPrimitive()) { + if (field.equals(targetField)) { + return new GroupType(candidate.getRepetition(), candidate.getName(), targetField); + } + } else { + Type tempField = extractField(field.asGroupType(), targetField); + if (tempField != null) { + return tempField; + } + } + } + + return null; + } + + private static final class DummyGroupConverter extends GroupConverter { + @Override + public void start() { + } + + @Override + public void end() { + } + + @Override + public Converter getConverter(int fieldIndex) { + return new DummyConverter(); + } + } + + private static final class DummyConverter extends PrimitiveConverter { + @Override + public GroupConverter asGroupConverter() { + return new DummyGroupConverter(); + } + } +} \ No newline at end of file diff --git a/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetRewriterFactory.java b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetRewriterFactory.java new file mode 100644 index 00000000000..61e5531fbaa --- /dev/null +++ b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetRewriterFactory.java @@ -0,0 +1,54 @@ +package org.apache.hudi.parquet.io; + +import org.apache.hudi.common.config.HoodieConfig; +import org.apache.hudi.common.config.HoodieStorageConfig; +import org.apache.hudi.common.util.ValidationUtils; +import org.apache.hudi.io.storage.HoodieParquetConfig; +import org.apache.hudi.io.storage.rewrite.HoodieFileMetadataMerger; +import org.apache.hudi.io.storage.rewrite.HoodieFileRewriter; +import org.apache.hudi.io.storage.rewrite.HoodieFileRewriterFactory; +import org.apache.hudi.storage.StoragePath; + +import org.apache.avro.Schema; +import org.apache.hadoop.conf.Configuration; +import org.apache.parquet.avro.AvroSchemaConverter; +import org.apache.parquet.hadoop.metadata.CompressionCodecName; +import org.apache.parquet.schema.MessageType; + +import java.io.IOException; +import java.util.List; + +public class HoodieParquetRewriterFactory extends HoodieFileRewriterFactory { + + @Override + protected <T> HoodieFileRewriter newFileRewriter( + List<StoragePath> inputFilePaths, + StoragePath targetFilePath, + Configuration conf, + HoodieConfig config, + HoodieFileMetadataMerger metadataMerger, + Schema writeSchemaWithMetaFields) throws IOException { + + ValidationUtils.checkArgument(writeSchemaWithMetaFields != null, + "write schema for ParquetFileRewriter can not be null"); + MessageType writeSchema = new AvroSchemaConverter(conf).convert(writeSchemaWithMetaFields); + + String compressionCodecName = config.getStringOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME); + HoodieParquetConfig parquetConfig = new HoodieParquetConfig(null, + CompressionCodecName.fromConf(compressionCodecName.isEmpty() ? null : compressionCodecName), + config.getIntOrDefault(HoodieStorageConfig.PARQUET_BLOCK_SIZE), + config.getIntOrDefault(HoodieStorageConfig.PARQUET_PAGE_SIZE), + config.getLongOrDefault(HoodieStorageConfig.PARQUET_MAX_FILE_SIZE), + null, + config.getDoubleOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_RATIO_FRACTION), + config.getBooleanOrDefault(HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED)); + + return new HoodieParquetFileRewriter( + inputFilePaths, + targetFilePath, + conf, + parquetConfig, + metadataMerger, + writeSchema); + } +} diff --git a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestFile.java b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestFile.java new file mode 100644 index 00000000000..92351e1a828 --- /dev/null +++ b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestFile.java @@ -0,0 +1,39 @@ +/* + * 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.hudi.parquet.io; + +import org.apache.parquet.example.data.simple.SimpleGroup; + +public class TestFile { + private final String fileName; + private final SimpleGroup[] fileContent; + + public TestFile(String fileName, SimpleGroup[] fileContent) { + this.fileName = fileName; + this.fileContent = fileContent; + } + + public String getFileName() { + return this.fileName; + } + + public SimpleGroup[] getFileContent() { + return this.fileContent; + } +} diff --git a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestFileBuilder.java b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestFileBuilder.java new file mode 100644 index 00000000000..62ae74ef6ee --- /dev/null +++ b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestFileBuilder.java @@ -0,0 +1,188 @@ +/* + * 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.hudi.parquet.io; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.parquet.column.ParquetProperties; +import org.apache.parquet.crypto.ParquetCipher; +import org.apache.parquet.example.data.Group; +import org.apache.parquet.example.data.simple.SimpleGroup; +import org.apache.parquet.hadoop.ParquetWriter; +import org.apache.parquet.hadoop.example.ExampleParquetWriter; +import org.apache.parquet.hadoop.example.GroupWriteSupport; +import org.apache.parquet.hadoop.metadata.CompressionCodecName; +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.Type; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Paths; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.ThreadLocalRandom; + +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.BINARY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT32; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT64; + +public class TestFileBuilder { + private MessageType schema; + private Configuration conf; + private Map<String, String> extraMeta = new HashMap<>(); + private int numRecord = 100000; + private ParquetProperties.WriterVersion writerVersion = ParquetProperties.WriterVersion.PARQUET_1_0; + private int pageSize = ParquetProperties.DEFAULT_PAGE_SIZE; + private String codec = "ZSTD"; + private String[] encryptColumns = {}; + private ParquetCipher cipher = ParquetCipher.AES_GCM_V1; + private Boolean footerEncryption = false; + + public TestFileBuilder(Configuration conf, MessageType schema) { + this.conf = conf; + this.schema = schema; + conf.set(GroupWriteSupport.PARQUET_EXAMPLE_SCHEMA, schema.toString()); + } + + public TestFileBuilder withNumRecord(int numRecord) { + this.numRecord = numRecord; + return this; + } + + public TestFileBuilder withEncrytionAlgorithm(ParquetCipher cipher) { + this.cipher = cipher; + return this; + } + + public TestFileBuilder withExtraMeta(Map<String, String> extraMeta) { + this.extraMeta = extraMeta; + return this; + } + + public TestFileBuilder withWriterVersion(ParquetProperties.WriterVersion writerVersion) { + this.writerVersion = writerVersion; + return this; + } + + public TestFileBuilder withPageSize(int pageSize) { + this.pageSize = pageSize; + return this; + } + + public TestFileBuilder withCodec(String codec) { + this.codec = codec; + return this; + } + + public TestFileBuilder withEncryptColumns(String[] encryptColumns) { + this.encryptColumns = encryptColumns; + return this; + } + + public TestFileBuilder withFooterEncryption() { + this.footerEncryption = true; + return this; + } + + public TestFile build() + throws IOException { + String fileName = createTempFile("test"); + SimpleGroup[] fileContent = createFileContent(schema); + ExampleParquetWriter.Builder builder = ExampleParquetWriter.builder(new Path(fileName)) + .withConf(conf) + .withWriterVersion(writerVersion) + .withExtraMetaData(extraMeta) + .withValidation(true) + .withPageSize(pageSize) + .withCompressionCodec(CompressionCodecName.valueOf(codec)); + try (ParquetWriter writer = builder.build()) { + for (int i = 0; i < fileContent.length; i++) { + writer.write(fileContent[i]); + } + } + return new TestFile(fileName, fileContent); + } + + private SimpleGroup[] createFileContent(MessageType schema) { + SimpleGroup[] simpleGroups = new SimpleGroup[numRecord]; + for (int i = 0; i < simpleGroups.length; i++) { + SimpleGroup g = new SimpleGroup(schema); + for (Type type : schema.getFields()) { + addValueToSimpleGroup(g, type); + } + simpleGroups[i] = g; + } + return simpleGroups; + } + + private void addValueToSimpleGroup(Group g, Type type) { + if (type.isPrimitive()) { + PrimitiveType primitiveType = (PrimitiveType) type; + if (primitiveType.getPrimitiveTypeName().equals(INT32)) { + g.add(type.getName(), getInt()); + } else if (primitiveType.getPrimitiveTypeName().equals(INT64)) { + g.add(type.getName(), getLong()); + } else if (primitiveType.getPrimitiveTypeName().equals(BINARY)) { + g.add(type.getName(), getString()); + } + // Only support 3 types now, more can be added later + } else { + GroupType groupType = (GroupType) type; + Group parentGroup = g.addGroup(groupType.getName()); + for (Type field : groupType.getFields()) { + addValueToSimpleGroup(parentGroup, field); + } + } + } + + private static long getInt() { + return ThreadLocalRandom.current().nextInt(10000); + } + + private static long getLong() { + return ThreadLocalRandom.current().nextLong(100000); + } + + private static String getString() { + char[] chars = {'a', 'b', 'c', 'd', 'e', 'f', 'g', 'x', 'z', 'y'}; + StringBuilder sb = new StringBuilder(); + for (int i = 0; i < 100; i++) { + sb.append(chars[ThreadLocalRandom.current().nextInt(10)]); + } + return sb.toString(); + } + + public static String createTempFile(String prefix) { + try { + return Files.createTempDirectory(prefix).toAbsolutePath().toString() + "/test.parquet"; + } catch (IOException e) { + throw new AssertionError("Unable to create temporary file", e); + } + } + + public static void deleteTempFile(String file) { + try { + Files.delete(Paths.get(file)); + } catch (IOException e) { + e.printStackTrace(); + } + } +} diff --git a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetFileRewriter.java b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetFileRewriter.java new file mode 100644 index 00000000000..ec477f1f838 --- /dev/null +++ b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetFileRewriter.java @@ -0,0 +1,475 @@ +/* + * 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.hudi.parquet.io; + +import org.apache.hudi.common.config.HoodieConfig; +import org.apache.hudi.common.config.HoodieStorageConfig; +import org.apache.hudi.io.storage.HoodieParquetConfig; +import org.apache.hudi.io.storage.rewrite.HoodieFileMetadataMerger; +import org.apache.hudi.storage.StoragePath; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.parquet.HadoopReadOptions; +import org.apache.parquet.ParquetReadOptions; +import org.apache.parquet.Version; +import org.apache.parquet.column.ParquetProperties; +import org.apache.parquet.example.data.Group; +import org.apache.parquet.example.data.simple.SimpleGroup; +import org.apache.parquet.format.DataPageHeader; +import org.apache.parquet.format.DataPageHeaderV2; +import org.apache.parquet.format.PageHeader; +import org.apache.parquet.format.converter.ParquetMetadataConverter; +import org.apache.parquet.hadoop.ParquetFileReader; +import org.apache.parquet.hadoop.ParquetReader; +import org.apache.parquet.hadoop.example.GroupReadSupport; +import org.apache.parquet.hadoop.metadata.BlockMetaData; +import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData; +import org.apache.parquet.hadoop.metadata.CompressionCodecName; +import org.apache.parquet.hadoop.metadata.FileMetaData; +import org.apache.parquet.hadoop.metadata.ParquetMetadata; +import org.apache.parquet.hadoop.util.CompressionConverter.TransParquetFileReader; +import org.apache.parquet.hadoop.util.HadoopInputFile; +import org.apache.parquet.internal.column.columnindex.ColumnIndex; +import org.apache.parquet.internal.column.columnindex.OffsetIndex; +import org.apache.parquet.io.InputFile; +import org.apache.parquet.io.SeekableInputStream; +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.Type; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; + +import static org.apache.hudi.common.model.HoodieRecord.FILENAME_METADATA_FIELD; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.BINARY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT64; +import static org.apache.parquet.schema.Type.Repetition.OPTIONAL; +import static org.apache.parquet.schema.Type.Repetition.REPEATED; +import static org.apache.parquet.schema.Type.Repetition.REQUIRED; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; + +public class TestHoodieParquetFileRewriter { + + private final int numRecord = 1; + private Configuration conf = new Configuration(); + private List<TestFile> inputFiles = null; + private String outputFile = null; + private HoodieParquetFileRewriter rewriter = null; + + @BeforeEach + public void setUp() { + outputFile = TestFileBuilder.createTempFile("test"); + } + + @AfterEach + public void after() { + if (outputFile != null) { + TestFileBuilder.deleteTempFile(outputFile); + } + if (inputFiles != null) { + inputFiles.stream().map(TestFile::getFileName).forEach(TestFileBuilder::deleteTempFile); + } + } + + @Test + public void testBasic() throws Exception { + MessageType schema = createSchema(); + inputFiles = new ArrayList<>(); + inputFiles.add(makeTestFile(schema, "GZIP")); + inputFiles.add(makeTestFile(schema, "GZIP")); + + rewriter = parquetFileRewriter(schema, "GZIP"); + rewriter.rewrite(); + rewriter.close(); + + // Verify the schema are not changed + ParquetMetadata pmd = ParquetFileReader.readFooter(conf, new Path(outputFile), ParquetMetadataConverter.NO_FILTER); + MessageType fileSchema = pmd.getFileMetaData().getSchema(); + assertEquals(schema, fileSchema); + validateSchema(fileSchema); + + // Verify codec + verifyCodec(outputFile, CompressionCodecName.GZIP); + + // Verify the merged data are not changed + validateColumnData(); + + // Verify the page index + validatePageIndex(0, 1, 2, 3, 4); + + // Verify original.created.by is preserved + validateCreatedBy(); + } + + @Test + public void testTranslateCodec() throws Exception { + MessageType schema = createSchema(); + inputFiles = new ArrayList<>(); + inputFiles.add(makeTestFile(schema, "GZIP")); + inputFiles.add(makeTestFile(schema, "UNCOMPRESSED")); + + rewriter = parquetFileRewriter(schema, "ZSTD"); + rewriter.rewrite(); + rewriter.close(); + + // Verify the schema are not changed for the columns not pruned + ParquetMetadata pmd = ParquetFileReader.readFooter(conf, new Path(outputFile), ParquetMetadataConverter.NO_FILTER); + MessageType fileSchema = pmd.getFileMetaData().getSchema(); + assertEquals(schema, fileSchema); + validateSchema(fileSchema); + + // Verify codec has been translated + verifyCodec(outputFile, CompressionCodecName.ZSTD); + + // Verify the data are not changed for the columns not pruned + validateColumnData(); + + // Verify the page index + validatePageIndex(0, 1, 2, 3, 4); + + // Verify original.created.by is preserved + validateCreatedBy(); + } + + @Test + public void testDifferentSchema() throws Exception { + MessageType schema1 = new MessageType("schema", + new PrimitiveType(OPTIONAL, INT64, "DocId"), + new PrimitiveType(REQUIRED, BINARY, "Name"), + new PrimitiveType(OPTIONAL, BINARY, "Gender"), + new GroupType(OPTIONAL, "Links", + new PrimitiveType(REPEATED, BINARY, "Backward"), + new PrimitiveType(REPEATED, BINARY, "Forward"))); + MessageType schema2 = new MessageType("schema", + new PrimitiveType(OPTIONAL, INT64, "DocId"), + new PrimitiveType(REQUIRED, BINARY, "Name"), + new PrimitiveType(OPTIONAL, BINARY, "Gender")); + inputFiles = new ArrayList<>(); + inputFiles.add(makeTestFile(schema1, "UNCOMPRESSED")); + inputFiles.add(makeTestFile(schema2, "UNCOMPRESSED")); + + rewriter = parquetFileRewriter(schema1, "UNCOMPRESSED"); + rewriter.rewrite(); + rewriter.close(); + + // Verify the schema are not changed for the columns not pruned + ParquetMetadata pmd = ParquetFileReader.readFooter(conf, new Path(outputFile), ParquetMetadataConverter.NO_FILTER); + MessageType schema = pmd.getFileMetaData().getSchema(); + validateSchema(schema); + + // Verify codec has been translated + verifyCodec(outputFile, CompressionCodecName.UNCOMPRESSED); + + // Verify the data are not changed + validateColumnData(); + + // Verify the page index + validatePageIndex(0, 1, 2); + + // Verify original.created.by is preserved + validateCreatedBy(); + } + + @Test + public void testHoodieMetaColumn() throws Exception { + MessageType schema = new MessageType("schema", + new PrimitiveType(OPTIONAL, BINARY, FILENAME_METADATA_FIELD), + new PrimitiveType(OPTIONAL, INT64, "DocId"), + new PrimitiveType(REQUIRED, BINARY, "Name"), + new PrimitiveType(OPTIONAL, BINARY, "Gender"), + new GroupType(OPTIONAL, "Links", + new PrimitiveType(REPEATED, BINARY, "Backward"), + new PrimitiveType(REPEATED, BINARY, "Forward"))); + inputFiles = new ArrayList<>(); + inputFiles.add(makeTestFile(schema, "GZIP")); + inputFiles.add(makeTestFile(schema, "GZIP")); + + rewriter = parquetFileRewriter(schema, "GZIP"); + rewriter.rewrite(); + rewriter.close(); + + // Verify the schema are not changed for the columns not pruned + ParquetMetadata pmd = ParquetFileReader.readFooter(conf, new Path(outputFile), ParquetMetadataConverter.NO_FILTER); + MessageType fileSchema = pmd.getFileMetaData().getSchema(); + assertEquals(schema, fileSchema); + + // Verify codec has been translated + verifyCodec(outputFile, CompressionCodecName.GZIP); + + // Verify the data are not changed + validateColumnData(); + + // Verify the page index + validatePageIndex(1, 2, 3, 4); + + // Verify original.created.by is preserved + validateCreatedBy(); + } + + private TestFile makeTestFile(MessageType schema, String codec) throws IOException { + return new TestFileBuilder(conf, schema) + .withNumRecord(numRecord) + .withCodec(codec) + .withPageSize(ParquetProperties.DEFAULT_PAGE_SIZE) + .build(); + } + + private HoodieParquetFileRewriter parquetFileRewriter(MessageType schema, String codec) { + List<StoragePath> inputPaths = inputFiles.stream() + .map(TestFile::getFileName) + .map(StoragePath::new) + .collect(Collectors.toList()); + StoragePath outputPath = new StoragePath(outputFile); + HoodieParquetConfig parquetConfig = new HoodieParquetConfig(null, + CompressionCodecName.fromConf(codec), + 120 * 1024 * 1024, + 1 * 1024 * 1024, + 120 * 1024 * 1024, + null, + 0.1, + HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED.defaultValue()); + return new HoodieParquetFileRewriter(inputPaths, outputPath, conf, parquetConfig, new HoodieFileMetadataMerger(), schema); + } + + private MessageType createSchema() { + return new MessageType("schema", + new PrimitiveType(OPTIONAL, INT64, "DocId"), + new PrimitiveType(REQUIRED, BINARY, "Name"), + new PrimitiveType(OPTIONAL, BINARY, "Gender"), + new GroupType(OPTIONAL, "Links", + new PrimitiveType(REPEATED, BINARY, "Backward"), + new PrimitiveType(REPEATED, BINARY, "Forward"))); + } + + private void validateSchema(MessageType schema) { + List<Type> fields = schema.getFields(); + assertEquals(fields.size(), 4); + assertEquals(fields.get(0).getName(), "DocId"); + assertEquals(fields.get(1).getName(), "Name"); + assertEquals(fields.get(2).getName(), "Gender"); + assertEquals(fields.get(3).getName(), "Links"); + List<Type> subFields = fields.get(3).asGroupType().getFields(); + assertEquals(subFields.size(), 2); + assertEquals(subFields.get(0).getName(), "Backward"); + assertEquals(subFields.get(1).getName(), "Forward"); + } + + private void validateColumnData() throws IOException { + Path outputFilePath = new Path(outputFile); + ParquetReader<Group> reader = ParquetReader.builder(new GroupReadSupport(), outputFilePath) + .withConf(conf) + .build(); + + // Get total number of rows from input files + int totalRows = 0; + for (TestFile inputFile : inputFiles) { + totalRows += inputFile.getFileContent().length; + } + + for (int i = 0; i < totalRows; i++) { + Group group = reader.read(); + assertNotNull(group); + + SimpleGroup expectGroup = inputFiles.get(i / numRecord).getFileContent()[i % numRecord]; + if (group.getType().containsField(FILENAME_METADATA_FIELD)) { + assertEquals(group.getString(FILENAME_METADATA_FIELD, 0), outputFilePath.getName()); + assertNotEquals(group.getString(FILENAME_METADATA_FIELD, 0), + expectGroup.getString(FILENAME_METADATA_FIELD, 0)); + } + assertEquals(group.getLong("DocId", 0), expectGroup.getLong("DocId", 0)); + assertArrayEquals(group.getBinary("Name", 0).getBytes(), + expectGroup.getBinary("Name", 0).getBytes()); + assertArrayEquals(group.getBinary("Gender", 0).getBytes(), + expectGroup.getBinary("Gender", 0).getBytes()); + + if (expectGroup.getType().containsField("Links")) { + Group subGroup = group.getGroup("Links", 0); + Group expectSubGroup = expectGroup.getGroup("Links", 0); + assertArrayEquals(subGroup.getBinary("Backward", 0).getBytes(), + expectSubGroup.getBinary("Backward", 0).getBytes()); + assertArrayEquals(subGroup.getBinary("Forward", 0).getBytes(), + expectSubGroup.getBinary("Forward", 0).getBytes()); + } + } + + reader.close(); + } + + private ParquetMetadata getFileMetaData(String file) throws IOException { + ParquetReadOptions readOptions = ParquetReadOptions.builder().build(); + InputFile inputFile = HadoopInputFile.fromPath(new Path(file), conf); + try (SeekableInputStream in = inputFile.newStream()) { + return ParquetFileReader.readFooter(inputFile, readOptions, in); + } + } + + private void verifyCodec(String file, CompressionCodecName expectedCodecs) throws IOException { + Set<CompressionCodecName> codecs = new HashSet<>(); + ParquetMetadata pmd = getFileMetaData(file); + for (int i = 0; i < pmd.getBlocks().size(); i++) { + BlockMetaData block = pmd.getBlocks().get(i); + for (int j = 0; j < block.getColumns().size(); ++j) { + ColumnChunkMetaData columnChunkMetaData = block.getColumns().get(j); + codecs.add(columnChunkMetaData.getCodec()); + } + } + assertEquals(new HashSet<CompressionCodecName>() { + { + add(expectedCodecs); + } + }, codecs); + } + + /** + * Verify the page index is correct. + * + * @param columnIdxs the idx of column to be validated. + */ + private void validatePageIndex(Integer... columnIdxs) throws Exception { + ParquetMetadata outMetaData = getFileMetaData(outputFile); + + int inputFileIndex = 0; + TransParquetFileReader inReader = new TransParquetFileReader( + HadoopInputFile.fromPath(new Path(inputFiles.get(inputFileIndex).getFileName()), conf), + HadoopReadOptions.builder(conf).build() + ); + ParquetMetadata inMetaData = inReader.getFooter(); + + try (TransParquetFileReader outReader = new TransParquetFileReader( + HadoopInputFile.fromPath(new Path(outputFile), conf), + HadoopReadOptions.builder(conf).build())) { + + for (int outBlockId = 0, inBlockId = 0; outBlockId < outMetaData.getBlocks().size(); ++outBlockId, ++inBlockId) { + // Refresh reader of input file + if (inBlockId == inMetaData.getBlocks().size()) { + inReader = new TransParquetFileReader( + HadoopInputFile.fromPath(new Path(inputFiles.get(++inputFileIndex).getFileName()), conf), + HadoopReadOptions.builder(conf).build()); + inMetaData = inReader.getFooter(); + inBlockId = 0; + } + + BlockMetaData inBlockMetaData = inMetaData.getBlocks().get(inBlockId); + BlockMetaData outBlockMetaData = outMetaData.getBlocks().get(outBlockId); + + for (int j = 0; j < columnIdxs.length; j++) { + ColumnChunkMetaData inChunk = inBlockMetaData.getColumns().get(columnIdxs[j]); + ColumnIndex inColumnIndex = inReader.readColumnIndex(inChunk); + OffsetIndex inOffsetIndex = inReader.readOffsetIndex(inChunk); + ColumnChunkMetaData outChunk = outBlockMetaData.getColumns().get(columnIdxs[j]); + ColumnIndex outColumnIndex = outReader.readColumnIndex(outChunk); + OffsetIndex outOffsetIndex = outReader.readOffsetIndex(outChunk); + if (inColumnIndex != null) { + assertEquals(inColumnIndex.getBoundaryOrder(), outColumnIndex.getBoundaryOrder()); + assertEquals(inColumnIndex.getMaxValues(), outColumnIndex.getMaxValues()); + assertEquals(inColumnIndex.getMinValues(), outColumnIndex.getMinValues()); + assertEquals(inColumnIndex.getNullCounts(), outColumnIndex.getNullCounts()); + } + if (inOffsetIndex != null) { + List<Long> inOffsets = getOffsets(inReader, inChunk); + List<Long> outOffsets = getOffsets(outReader, outChunk); + assertEquals(inOffsets.size(), outOffsets.size()); + assertEquals(inOffsets.size(), inOffsetIndex.getPageCount()); + assertEquals(inOffsetIndex.getPageCount(), outOffsetIndex.getPageCount()); + for (int k = 0; k < inOffsetIndex.getPageCount(); k++) { + assertEquals(inOffsetIndex.getFirstRowIndex(k), outOffsetIndex.getFirstRowIndex(k)); + assertEquals(inOffsetIndex.getLastRowIndex(k, inChunk.getValueCount()), + outOffsetIndex.getLastRowIndex(k, outChunk.getValueCount())); + assertEquals(inOffsetIndex.getOffset(k), (long) inOffsets.get(k)); + assertEquals(outOffsetIndex.getOffset(k), (long) outOffsets.get(k)); + } + } + } + } + } + } + + private List<Long> getOffsets(TransParquetFileReader reader, ColumnChunkMetaData chunk) throws IOException { + List<Long> offsets = new ArrayList<>(); + reader.setStreamPosition(chunk.getStartingPos()); + long readValues = 0; + long totalChunkValues = chunk.getValueCount(); + while (readValues < totalChunkValues) { + long curOffset = reader.getPos(); + PageHeader pageHeader = reader.readPageHeader(); + switch (pageHeader.type) { + case DICTIONARY_PAGE: + rewriter.readBlock(pageHeader.getCompressed_page_size(), reader); + break; + case DATA_PAGE: + DataPageHeader headerV1 = pageHeader.data_page_header; + offsets.add(curOffset); + rewriter.readBlock(pageHeader.getCompressed_page_size(), reader); + readValues += headerV1.getNum_values(); + break; + case DATA_PAGE_V2: + DataPageHeaderV2 headerV2 = pageHeader.data_page_header_v2; + offsets.add(curOffset); + int rlLength = headerV2.getRepetition_levels_byte_length(); + rewriter.readBlock(rlLength, reader); + int dlLength = headerV2.getDefinition_levels_byte_length(); + rewriter.readBlock(dlLength, reader); + int payLoadLength = pageHeader.getCompressed_page_size() - rlLength - dlLength; + rewriter.readBlock(payLoadLength, reader); + readValues += headerV2.getNum_values(); + break; + default: + throw new IOException("Not recognized page type"); + } + } + return offsets; + } + + private void validateCreatedBy() throws Exception { + Set<String> createdBySet = new HashSet<>(); + for (TestFile inputFile : inputFiles) { + ParquetMetadata pmd = getFileMetaData(inputFile.getFileName()); + createdBySet.add(pmd.getFileMetaData().getCreatedBy()); + assertNull(pmd.getFileMetaData().getKeyValueMetaData().get(HoodieParquetFileRewriter.ORIGINAL_CREATED_BY_KEY)); + } + + // Verify created_by from input files have been deduplicated + Object[] inputCreatedBys = createdBySet.toArray(); + assertEquals(1, inputCreatedBys.length); + + // Verify created_by has been set + FileMetaData outFMD = getFileMetaData(outputFile).getFileMetaData(); + final String createdBy = outFMD.getCreatedBy(); + assertNotNull(createdBy); + assertEquals(createdBy, Version.FULL_VERSION); + + // Verify original.created.by has been set + String inputCreatedBy = (String) inputCreatedBys[0]; + String originalCreatedBy = outFMD.getKeyValueMetaData().get(HoodieParquetFileRewriter.ORIGINAL_CREATED_BY_KEY); + assertEquals(inputCreatedBy, originalCreatedBy); + } +} \ No newline at end of file
