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

Reply via email to