This is an automated email from the ASF dual-hosted git repository.

danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 3c680be7d278 [HUDI-9704] Move remaining APIs from reader context to 
record context (#13717)
3c680be7d278 is described below

commit 3c680be7d278ceb5e630d246cca6641ff7c0696e
Author: Lokesh Jain <[email protected]>
AuthorDate: Wed Aug 13 19:17:32 2025 +0530

    [HUDI-9704] Move remaining APIs from reader context to record context 
(#13717)
    
    * [HUDI-9704] Move remaining APIs from reader context to record context
    * Fix compilation
    
    ---------
    
    Co-authored-by: Lokesh Jain <[email protected]>
---
 .../hudi/BaseSparkInternalRecordContext.java       | 27 +++++++++++
 .../hudi/BaseSparkInternalRowReaderContext.java    | 22 ---------
 .../SparkFileFormatInternalRowReaderContext.scala  |  2 +-
 .../TestBaseSparkInternalRowReaderContext.java     | 10 ++---
 .../org/apache/hudi/avro/AvroRecordContext.java    | 16 +++++++
 .../apache/hudi/avro/HoodieAvroReaderContext.java  | 16 -------
 .../hudi/common/engine/HoodieReaderContext.java    | 37 ---------------
 .../apache/hudi/common/engine/RecordContext.java   | 36 +++++++++++++++
 .../hudi/common/table/read/BufferedRecord.java     | 10 ++---
 .../table/read/FileGroupReaderSchemaHandler.java   |  2 +-
 .../common/table/read/HoodieFileGroupReader.java   |  2 +-
 .../hudi/common/table/read/UpdateProcessor.java    |  2 +-
 .../table/read/buffer/FileGroupRecordBuffer.java   |  4 +-
 .../read/buffer/KeyBasedFileGroupRecordBuffer.java |  2 +-
 .../read/buffer/UnmergedFileGroupRecordBuffer.java |  4 +-
 .../common/table/read/SchemaHandlerTestBase.java   | 30 ++++++-------
 .../table/read/TestHoodieFileGroupReaderBase.java  |  4 +-
 .../buffer/TestReusableKeyBasedRecordBuffer.java   |  4 +-
 .../TestSortedKeyBasedFileGroupRecordBuffer.java   |  4 +-
 .../hudi/table/format/FlinkRecordContext.java      | 42 +++++++++++++++++
 .../table/format/FlinkRowDataReaderContext.java    | 52 +++-------------------
 .../hudi/hadoop/HiveHoodieReaderContext.java       | 21 +--------
 .../org/apache/hudi/hadoop/HiveRecordContext.java  | 18 ++++++++
 .../hudi/functional/TestBufferedRecordMerger.java  | 32 ++++++-------
 24 files changed, 201 insertions(+), 198 deletions(-)

diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
index a79f65a99776..5ea36a6e82de 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRecordContext.java
@@ -35,11 +35,17 @@ import org.apache.spark.sql.HoodieInternalRowUtils;
 import org.apache.spark.sql.HoodieUnsafeRowUtils;
 import org.apache.spark.sql.catalyst.InternalRow;
 import org.apache.spark.sql.catalyst.expressions.GenericInternalRow;
+import org.apache.spark.sql.catalyst.expressions.UnsafeProjection;
+import org.apache.spark.sql.catalyst.expressions.UnsafeRow;
 import org.apache.spark.sql.types.StructType;
 import org.apache.spark.unsafe.types.UTF8String;
 
+import java.util.Collections;
 import java.util.List;
 import java.util.Map;
+import java.util.function.UnaryOperator;
+
+import scala.Function1;
 
 import static org.apache.spark.sql.HoodieInternalRowUtils.getCachedSchema;
 
@@ -120,4 +126,25 @@ public abstract class BaseSparkInternalRecordContext 
extends RecordContext<Inter
   public InternalRow getDeleteRow(String recordKey) {
     return new HoodieInternalRow(null, null, UTF8String.fromString(recordKey), 
UTF8String.fromString(partitionPath), null, null, false);
   }
+
+  @Override
+  public InternalRow seal(InternalRow internalRow) {
+    return internalRow.copy();
+  }
+
+  @Override
+  public InternalRow toBinaryRow(Schema schema, InternalRow internalRow) {
+    if (internalRow instanceof UnsafeRow) {
+      return internalRow;
+    }
+    final UnsafeProjection unsafeProjection = 
HoodieInternalRowUtils.getCachedUnsafeProjection(schema);
+    return unsafeProjection.apply(internalRow);
+  }
+
+  @Override
+  public UnaryOperator<InternalRow> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
+    Function1<InternalRow, UnsafeRow> unsafeRowWriter =
+        HoodieInternalRowUtils.getCachedUnsafeRowWriter(getCachedSchema(from), 
getCachedSchema(to), renamedColumns, Collections.emptyMap());
+    return row -> (InternalRow) unsafeRowWriter.apply(row);
+  }
 }
diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
index 0c581003a9b0..d8d384352cd2 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
@@ -32,7 +32,6 @@ import org.apache.hudi.storage.StorageConfiguration;
 import org.apache.avro.Schema;
 import org.apache.spark.sql.HoodieInternalRowUtils;
 import org.apache.spark.sql.catalyst.InternalRow;
-import org.apache.spark.sql.catalyst.expressions.UnsafeProjection;
 import org.apache.spark.sql.catalyst.expressions.UnsafeRow;
 
 import java.util.Collections;
@@ -78,27 +77,6 @@ public abstract class BaseSparkInternalRowReaderContext 
extends HoodieReaderCont
     }
   }
 
-  @Override
-  public InternalRow seal(InternalRow internalRow) {
-    return internalRow.copy();
-  }
-
-  @Override
-  public InternalRow toBinaryRow(Schema schema, InternalRow internalRow) {
-    if (internalRow instanceof UnsafeRow) {
-      return internalRow;
-    }
-    final UnsafeProjection unsafeProjection = 
HoodieInternalRowUtils.getCachedUnsafeProjection(schema);
-    return unsafeProjection.apply(internalRow);
-  }
-
-  @Override
-  public UnaryOperator<InternalRow> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
-    Function1<InternalRow, UnsafeRow> unsafeRowWriter =
-        HoodieInternalRowUtils.getCachedUnsafeRowWriter(getCachedSchema(from), 
getCachedSchema(to), renamedColumns, Collections.emptyMap());
-    return row -> (InternalRow) unsafeRowWriter.apply(row);
-  }
-
   /**
    * Constructs a transformation that will take a row and convert it to a new 
row with the given schema and adds in the values for the partition columns if 
they are missing in the returned row.
    * It is assumed that the `to` schema will contain the partition fields.
diff --git 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
index aec62834efde..dfc4655576e9 100644
--- 
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
+++ 
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
@@ -130,7 +130,7 @@ class 
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
       val rowIndexColumn = new java.util.HashSet[String]()
       rowIndexColumn.add(ROW_INDEX_TEMPORARY_COLUMN_NAME)
       //always remove the row index column from the skeleton because the data 
file will also have the same column
-      val skeletonProjection = projectRecord(skeletonRequiredSchema,
+      val skeletonProjection = 
recordContext.projectRecord(skeletonRequiredSchema,
         HoodieAvroUtils.removeFields(skeletonRequiredSchema, rowIndexColumn))
 
       //If we need to do position based merging with log files we will leave 
the row index column at the end
diff --git 
a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/TestBaseSparkInternalRowReaderContext.java
 
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/TestBaseSparkInternalRowReaderContext.java
index 2b0d299647a2..11287654d9c5 100644
--- 
a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/TestBaseSparkInternalRowReaderContext.java
+++ 
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/TestBaseSparkInternalRowReaderContext.java
@@ -163,12 +163,12 @@ class TestBaseSparkInternalRowReaderContext {
             return row.getBoolean(2);
           }
         }
-      });
-    }
 
-    @Override
-    public InternalRow toBinaryRow(Schema schema, InternalRow internalRow) {
-      return internalRow;
+        @Override
+        public InternalRow toBinaryRow(Schema schema, InternalRow internalRow) 
{
+          return internalRow;
+        }
+      });
     }
 
     @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/avro/AvroRecordContext.java 
b/hudi-common/src/main/java/org/apache/hudi/avro/AvroRecordContext.java
index 0767995cf5f6..a30fc6fd987d 100644
--- a/hudi-common/src/main/java/org/apache/hudi/avro/AvroRecordContext.java
+++ b/hudi-common/src/main/java/org/apache/hudi/avro/AvroRecordContext.java
@@ -41,6 +41,7 @@ import org.apache.avro.generic.IndexedRecord;
 import java.io.IOException;
 import java.util.Map;
 import java.util.Properties;
+import java.util.function.UnaryOperator;
 
 /**
  * Record context for reading and transforming avro indexed records.
@@ -165,4 +166,19 @@ public class AvroRecordContext extends 
RecordContext<IndexedRecord> {
   public IndexedRecord getDeleteRow(String recordKey) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
+
+  @Override
+  public IndexedRecord seal(IndexedRecord record) {
+    return record;
+  }
+
+  @Override
+  public IndexedRecord toBinaryRow(Schema avroSchema, IndexedRecord record) {
+    return record;
+  }
+
+  @Override
+  public UnaryOperator<IndexedRecord> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
+    return record -> HoodieAvroUtils.rewriteRecordWithNewSchema(record, to, 
renamedColumns);
+  }
 }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroReaderContext.java 
b/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroReaderContext.java
index f376b6c1a7c3..d966e9ffa788 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroReaderContext.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroReaderContext.java
@@ -53,7 +53,6 @@ import java.io.IOException;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
-import java.util.function.UnaryOperator;
 
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY;
 import static org.apache.hudi.common.util.ValidationUtils.checkState;
@@ -179,16 +178,6 @@ public class HoodieAvroReaderContext extends 
HoodieReaderContext<IndexedRecord>
     }
   }
 
-  @Override
-  public IndexedRecord seal(IndexedRecord record) {
-    return record;
-  }
-
-  @Override
-  public IndexedRecord toBinaryRow(Schema avroSchema, IndexedRecord record) {
-    return record;
-  }
-
   @Override
   public SizeEstimator<BufferedRecord<IndexedRecord>> getRecordSizeEstimator() 
{
     return new AvroRecordSizeEstimator(getSchemaHandler().getRequiredSchema());
@@ -208,11 +197,6 @@ public class HoodieAvroReaderContext extends 
HoodieReaderContext<IndexedRecord>
     return new BootstrapIterator(skeletonFileIterator, skeletonRequiredSchema, 
dataFileIterator, dataRequiredSchema, partitionFieldAndValues);
   }
 
-  @Override
-  public UnaryOperator<IndexedRecord> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
-    return record -> HoodieAvroUtils.rewriteRecordWithNewSchema(record, to, 
renamedColumns);
-  }
-
   /**
    * Iterator that traverses the skeleton file and the base file in tandem.
    * The iterator will only extract the fields requested in the provided 
schemas.
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
index 120b7cef9fdf..c498b836c396 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
@@ -51,10 +51,7 @@ import org.apache.hudi.storage.StoragePathInfo;
 import org.apache.avro.Schema;
 
 import java.io.IOException;
-import java.util.Collections;
 import java.util.List;
-import java.util.Map;
-import java.util.function.UnaryOperator;
 
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_DEPRECATED_WRITE_CONFIG_KEY;
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY;
@@ -321,24 +318,6 @@ public abstract class HoodieReaderContext<T> {
     return new CloseableFilterIterator<>(fileRecordIterator, instantFilter);
   }
 
-  /**
-   * Seals the engine-specific record to make sure the data referenced in 
memory do not change.
-   *
-   * @param record The record.
-   * @return The record containing the same data that do not change in memory 
over time.
-   */
-  public abstract T seal(T record);
-
-  /**
-   * Converts engine specific row into binary format.
-   *
-   * @param avroSchema The avro schema of the row
-   * @param record     The engine row
-   *
-   * @return row in binary format
-   */
-  public abstract T toBinaryRow(Schema avroSchema, T record);
-
   /**
    * Merge the skeleton file and data file iterators into a single iterator 
that will produce rows that contain all columns from the
    * skeleton file iterator, followed by all columns in the data file iterator
@@ -355,20 +334,4 @@ public abstract class HoodieReaderContext<T> {
                                                             
ClosableIterator<T> dataFileIterator,
                                                             Schema 
dataRequiredSchema,
                                                             List<Pair<String, 
Object>> requiredPartitionFieldAndValues);
-
-  /**
-   * Creates a function that will reorder records of schema "from" to schema 
of "to"
-   * all fields in "to" must be in "from", but not all fields in "from" must 
be in "to"
-   *
-   * @param from           the schema of records to be passed into 
UnaryOperator
-   * @param to             the schema of records produced by UnaryOperator
-   * @param renamedColumns map of renamed columns where the key is the new 
name from the query and
-   *                       the value is the old name that exists in the file
-   * @return a function that takes in a record and returns the record with 
reordered columns
-   */
-  public abstract UnaryOperator<T> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns);
-
-  public final UnaryOperator<T> projectRecord(Schema from, Schema to) {
-    return projectRecord(from, to, Collections.emptyMap());
-  }
 }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java 
b/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java
index f926c0e7c6f1..9c772971f119 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/engine/RecordContext.java
@@ -39,10 +39,12 @@ import org.apache.avro.generic.IndexedRecord;
 import javax.annotation.Nullable;
 
 import java.io.Serializable;
+import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.Properties;
 import java.util.function.BiFunction;
+import java.util.function.UnaryOperator;
 
 import static 
org.apache.hudi.common.model.HoodieRecord.HOODIE_IS_DELETED_FIELD;
 import static 
org.apache.hudi.common.model.HoodieRecord.RECORD_KEY_METADATA_FIELD;
@@ -274,6 +276,40 @@ public abstract class RecordContext<T> implements 
Serializable {
         && markerKeyValue.getRight().equals(deleteMarkerValue.toString());
   }
 
+  /**
+   * Seals the engine-specific record to make sure the data referenced in 
memory do not change.
+   *
+   * @param record The record.
+   * @return The record containing the same data that do not change in memory 
over time.
+   */
+  public abstract T seal(T record);
+
+  /**
+   * Converts engine specific row into binary format.
+   *
+   * @param avroSchema The avro schema of the row
+   * @param record     The engine row
+   *
+   * @return row in binary format
+   */
+  public abstract T toBinaryRow(Schema avroSchema, T record);
+
+  /**
+   * Creates a function that will reorder records of schema "from" to schema 
of "to"
+   * all fields in "to" must be in "from", but not all fields in "from" must 
be in "to"
+   *
+   * @param from           the schema of records to be passed into 
UnaryOperator
+   * @param to             the schema of records produced by UnaryOperator
+   * @param renamedColumns map of renamed columns where the key is the new 
name from the query and
+   *                       the value is the old name that exists in the file
+   * @return a function that takes in a record and returns the record with 
reordered columns
+   */
+  public abstract UnaryOperator<T> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns);
+
+  public final UnaryOperator<T> projectRecord(Schema from, Schema to) {
+    return projectRecord(from, to, Collections.emptyMap());
+  }
+
   /**
    * Gets the ordering value in particular type.
    *
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java
index dd57bdbc0496..06e22989810f 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java
@@ -18,7 +18,7 @@
 
 package org.apache.hudi.common.table.read;
 
-import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.engine.RecordContext;
 import org.apache.hudi.common.model.HoodieOperation;
 import org.apache.hudi.common.util.OrderingValues;
 
@@ -88,16 +88,16 @@ public class BufferedRecord<T> implements Serializable {
     return this.hoodieOperation;
   }
 
-  public BufferedRecord<T> toBinary(HoodieReaderContext<T> readerContext) {
+  public BufferedRecord<T> toBinary(RecordContext<T> recordContext) {
     if (record != null) {
-      record = 
readerContext.seal(readerContext.toBinaryRow(readerContext.getRecordContext().getSchemaFromBufferRecord(this),
 record));
+      record = 
recordContext.seal(recordContext.toBinaryRow(recordContext.getSchemaFromBufferRecord(this),
 record));
     }
     return this;
   }
 
-  public BufferedRecord<T> seal(HoodieReaderContext<T> readerContext) {
+  public BufferedRecord<T> seal(RecordContext<T> recordContext) {
     if (record != null) {
-      this.record = readerContext.seal(record);
+      this.record = recordContext.seal(record);
     }
     return this;
   }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
index a05fa38b2965..727077f4a52c 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
@@ -116,7 +116,7 @@ public class FileGroupReaderSchemaHandler<T> {
 
   public Option<UnaryOperator<T>> getOutputConverter() {
     if (!AvroSchemaUtils.areSchemasProjectionEquivalent(requiredSchema, 
requestedSchema)) {
-      return Option.of(readerContext.projectRecord(requiredSchema, 
requestedSchema));
+      return 
Option.of(readerContext.getRecordContext().projectRecord(requiredSchema, 
requestedSchema));
     }
     return Option.empty();
   }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
index 5615374c1313..54b8d2f63152 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
@@ -129,7 +129,7 @@ public final class HoodieFileGroupReader<T> implements 
Closeable {
   private void initRecordIterators() throws IOException {
     ClosableIterator<T> iter = makeBaseFileIterator();
     if (inputSplit.hasNoRecordsToMerge()) {
-      this.baseFileIterator = new CloseableMappingIterator<>(iter, 
readerContext::seal);
+      this.baseFileIterator = new CloseableMappingIterator<>(iter, rec -> 
readerContext.getRecordContext().seal(rec));
     } else {
       this.baseFileIterator = iter;
       Pair<HoodieFileGroupRecordBuffer<T>, List<String>> initializationResult 
= recordBufferLoader.getRecordBuffer(
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
index 4b095adb88ad..60a4151a11e0 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java
@@ -89,7 +89,7 @@ public interface UpdateProcessor<T> {
           mergedRecord.setHoodieOperation(HoodieOperation.INSERT);
           readStats.incrementNumInserts();
         }
-        return mergedRecord.seal(readerContext);
+        return mergedRecord.seal(readerContext.getRecordContext());
       }
     }
   }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/FileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/FileGroupRecordBuffer.java
index 53b60ba362a9..27c7be2d6831 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/FileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/FileGroupRecordBuffer.java
@@ -234,7 +234,7 @@ abstract class FileGroupRecordBuffer<T> implements 
HoodieFileGroupRecordBuffer<T
     // InternalSchemaMerger#buildRecordType() for details.
     // Delete and add a field with the same name, reads should not return 
previously inserted datum of dropped field of the same name,
     // so we use `mergedAvroSchema` as the target schema for record projecting.
-    return 
Option.of(Pair.of(readerContext.projectRecord(dataBlock.getSchema(), 
mergedAvroSchema, mergedInternalSchema.getRight()), mergedAvroSchema));
+    return 
Option.of(Pair.of(readerContext.getRecordContext().projectRecord(dataBlock.getSchema(),
 mergedAvroSchema, mergedInternalSchema.getRight()), mergedAvroSchema));
   }
 
   protected boolean hasNextBaseRecord(T baseRecord, BufferedRecord<T> 
logRecordInfo) throws IOException {
@@ -246,7 +246,7 @@ abstract class FileGroupRecordBuffer<T> implements 
HoodieFileGroupRecordBuffer<T
     }
 
     // Inserts
-    nextRecord = 
bufferedRecordConverter.convert(readerContext.seal(baseRecord));
+    nextRecord = 
bufferedRecordConverter.convert(readerContext.getRecordContext().seal(baseRecord));
     nextRecord.setHoodieOperation(HoodieOperation.INSERT);
     return true;
   }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/KeyBasedFileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/KeyBasedFileGroupRecordBuffer.java
index 8fe84cbc2042..f66629aaf519 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/KeyBasedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/KeyBasedFileGroupRecordBuffer.java
@@ -108,7 +108,7 @@ public class KeyBasedFileGroupRecordBuffer<T> extends 
FileGroupRecordBuffer<T> {
     BufferedRecord<T> existingRecord = records.get(recordKey);
     totalLogRecords++;
     bufferedRecordMerger.deltaMerge(record, 
existingRecord).ifPresent(bufferedRecord ->
-        records.put(recordKey, bufferedRecord.toBinary(readerContext)));
+        records.put(recordKey, 
bufferedRecord.toBinary(readerContext.getRecordContext())));
   }
 
   @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/UnmergedFileGroupRecordBuffer.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/UnmergedFileGroupRecordBuffer.java
index a4fc64b13e21..3574ef35552e 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/UnmergedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/buffer/UnmergedFileGroupRecordBuffer.java
@@ -67,7 +67,7 @@ class UnmergedFileGroupRecordBuffer<T> extends 
FileGroupRecordBuffer<T> {
 
     // Output from base file first.
     if (baseFileIterator.hasNext()) {
-      nextRecord = 
bufferedRecordConverter.convert(readerContext.seal(baseFileIterator.next()));
+      nextRecord = 
bufferedRecordConverter.convert(readerContext.getRecordContext().seal(baseFileIterator.next()));
       return true;
     }
 
@@ -85,7 +85,7 @@ class UnmergedFileGroupRecordBuffer<T> extends 
FileGroupRecordBuffer<T> {
     if (recordIterator == null || !recordIterator.hasNext()) {
       return false;
     }
-    nextRecord = 
bufferedRecordConverter.convert(readerContext.seal(recordIterator.next()));
+    nextRecord = 
bufferedRecordConverter.convert(readerContext.getRecordContext().seal(recordIterator.next()));
     readStats.incrementNumInserts();
     return true;
   }
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/SchemaHandlerTestBase.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/SchemaHandlerTestBase.java
index 41d63434f5d7..85813f59f1bc 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/SchemaHandlerTestBase.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/SchemaHandlerTestBase.java
@@ -346,6 +346,21 @@ public abstract class SchemaHandlerTestBase {
         public String mergeWithEngineRecord(Schema schema, Map<Integer, 
Object> updateValues, BufferedRecord<String> baseRecord) {
           return "";
         }
+
+        @Override
+        public String seal(String record) {
+          return "";
+        }
+
+        @Override
+        public String toBinaryRow(Schema avroSchema, String record) {
+          return "";
+        }
+
+        @Override
+        public UnaryOperator<String> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
+          return null;
+        }
       });
     }
 
@@ -359,25 +374,10 @@ public abstract class SchemaHandlerTestBase {
       return null;
     }
 
-    @Override
-    public String seal(String record) {
-      return "";
-    }
-
-    @Override
-    public String toBinaryRow(Schema avroSchema, String record) {
-      return "";
-    }
-
     @Override
     public ClosableIterator<String> 
mergeBootstrapReaders(ClosableIterator<String> skeletonFileIterator, Schema 
skeletonRequiredSchema, ClosableIterator<String> dataFileIterator,
                                                           Schema 
dataRequiredSchema, List<Pair<String, Object>> requiredPartitionFieldAndValues) 
{
       return null;
     }
-
-    @Override
-    public UnaryOperator<String> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
-      return null;
-    }
   }
 }
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
index df9e7c4793a5..a91e78afa584 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
@@ -596,10 +596,10 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
             String recordKey = 
readerContext.getRecordContext().getRecordKey(record, avroSchema);
             //test key based
             BufferedRecord<T> bufferedRecord = 
BufferedRecords.fromEngineRecord(record, avroSchema, 
readerContext.getRecordContext(), Collections.singletonList("timestamp"), 
false);
-            spillableMap.put(recordKey, 
bufferedRecord.toBinary(readerContext));
+            spillableMap.put(recordKey, 
bufferedRecord.toBinary(readerContext.getRecordContext()));
 
             //test position based
-            spillableMap.put(position++, 
bufferedRecord.toBinary(readerContext));
+            spillableMap.put(position++, 
bufferedRecord.toBinary(readerContext.getRecordContext()));
           }
 
           assertEquals(records.size() * 2, spillableMap.size());
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestReusableKeyBasedRecordBuffer.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestReusableKeyBasedRecordBuffer.java
index 5940e54422bc..6865899166af 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestReusableKeyBasedRecordBuffer.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestReusableKeyBasedRecordBuffer.java
@@ -85,8 +85,8 @@ class TestReusableKeyBasedRecordBuffer {
       // otherwise return an older value
       return 1;
     });
-    when(mockReaderContext.toBinaryRow(any(), any())).thenAnswer(invocation -> 
invocation.getArgument(1));
-    when(mockReaderContext.seal(any())).thenAnswer(invocation -> 
invocation.getArgument(0));
+    when(mockReaderContext.getRecordContext().toBinaryRow(any(), 
any())).thenAnswer(invocation -> invocation.getArgument(1));
+    
when(mockReaderContext.getRecordContext().seal(any())).thenAnswer(invocation -> 
invocation.getArgument(0));
 
     ReusableKeyBasedRecordBuffer<TestRecord> buffer = new 
ReusableKeyBasedRecordBuffer<>(mockReaderContext, metaClient,
         RecordMergeMode.EVENT_TIME_ORDERING, PartialUpdateMode.NONE, new 
TypedProperties(), Collections.singletonList("value"), updateProcessor, 
preMergedLogRecords);
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestSortedKeyBasedFileGroupRecordBuffer.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestSortedKeyBasedFileGroupRecordBuffer.java
index 0cbf8c2017e5..458979f1dfe0 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestSortedKeyBasedFileGroupRecordBuffer.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestSortedKeyBasedFileGroupRecordBuffer.java
@@ -195,8 +195,8 @@ class TestSortedKeyBasedFileGroupRecordBuffer extends 
BaseTestFileGroupRecordBuf
     });
     when(mockReaderContext.getRecordContext().getRecordKey(any(), 
any())).thenAnswer(invocation -> ((TestRecord) 
invocation.getArgument(0)).getRecordKey());
     when(mockReaderContext.getRecordContext().getOrderingValue(any(), any(), 
anyList())).thenReturn(0);
-    when(mockReaderContext.toBinaryRow(any(), any())).thenAnswer(invocation -> 
invocation.getArgument(1));
-    when(mockReaderContext.seal(any())).thenAnswer(invocation -> 
invocation.getArgument(0));
+    when(mockReaderContext.getRecordContext().toBinaryRow(any(), 
any())).thenAnswer(invocation -> invocation.getArgument(1));
+    
when(mockReaderContext.getRecordContext().seal(any())).thenAnswer(invocation -> 
invocation.getArgument(0));
     HoodieTableMetaClient mockMetaClient = mock(HoodieTableMetaClient.class);
     RecordMergeMode recordMergeMode = RecordMergeMode.COMMIT_TIME_ORDERING;
     PartialUpdateMode partialUpdateMode = PartialUpdateMode.NONE;
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRecordContext.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRecordContext.java
index 814f8eabeb9f..60a44f5bbfc6 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRecordContext.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRecordContext.java
@@ -35,16 +35,22 @@ import org.apache.hudi.util.AvroToRowDataConverters;
 import org.apache.hudi.util.RecordKeyToRowDataConverter;
 import org.apache.hudi.util.RowDataAvroQueryContexts;
 import org.apache.hudi.util.RowDataUtils;
+import org.apache.hudi.util.RowProjection;
+import org.apache.hudi.util.SchemaEvolvingRowDataProjection;
 
 import org.apache.avro.Schema;
 import org.apache.avro.generic.GenericRecord;
 import org.apache.avro.generic.IndexedRecord;
 import org.apache.flink.table.data.GenericRowData;
 import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.data.binary.BinaryRowData;
+import org.apache.flink.table.runtime.typeutils.RowDataSerializer;
+import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.types.RowKind;
 
 import java.util.List;
 import java.util.Map;
+import java.util.function.UnaryOperator;
 
 public class FlinkRecordContext extends RecordContext<RowData> {
 
@@ -156,6 +162,42 @@ public class FlinkRecordContext extends 
RecordContext<RowData> {
     });
   }
 
+  @Override
+  public RowData seal(RowData rowData) {
+    if (rowData instanceof BinaryRowData) {
+      return ((BinaryRowData) rowData).copy();
+    }
+    return rowData;
+  }
+
+  @Override
+  public RowData toBinaryRow(Schema avroSchema, RowData record) {
+    if (record instanceof BinaryRowData) {
+      return record;
+    }
+    RowDataSerializer rowDataSerializer = 
RowDataAvroQueryContexts.getRowDataSerializer(avroSchema);
+    return rowDataSerializer.toBinaryRow(record);
+  }
+
+  /**
+   * Creates a function that will reorder records of schema "from" to schema 
of "to".
+   * It's possible there exist fields in `to` schema, but not in `from` schema 
because of schema
+   * evolution.
+   *
+   * @param from           the schema of records to be passed into 
UnaryOperator
+   * @param to             the schema of records produced by UnaryOperator
+   * @param renamedColumns map of renamed columns where the key is the new 
name from the query and
+   *                       the value is the old name that exists in the file
+   * @return a function that takes in a record and returns the record with 
reordered columns
+   */
+  @Override
+  public UnaryOperator<RowData> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
+    RowType fromType = (RowType) 
RowDataAvroQueryContexts.fromAvroSchema(from).getRowType().getLogicalType();
+    RowType toType =  (RowType) 
RowDataAvroQueryContexts.fromAvroSchema(to).getRowType().getLogicalType();
+    RowProjection rowProjection = 
SchemaEvolvingRowDataProjection.instance(fromType, toType, renamedColumns);
+    return rowProjection::project;
+  }
+
   public void setRecordKeyRowConverter(RecordKeyToRowDataConverter 
recordKeyRowConverter) {
     this.recordKeyRowConverter = recordKeyRowConverter;
   }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
index 5d9dca452fa9..f10add8b0e9e 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
@@ -43,15 +43,11 @@ import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.StorageConfiguration;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.util.RowDataAvroQueryContexts;
-import org.apache.hudi.util.RowProjection;
-import org.apache.hudi.util.SchemaEvolvingRowDataProjection;
 import org.apache.hudi.util.RecordKeyToRowDataConverter;
 
 import org.apache.avro.Schema;
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.data.binary.BinaryRowData;
 import org.apache.flink.table.data.utils.JoinedRowData;
-import org.apache.flink.table.runtime.typeutils.RowDataSerializer;
 import org.apache.flink.table.types.DataType;
 import org.apache.flink.table.types.logical.RowType;
 
@@ -60,7 +56,6 @@ import java.util.Arrays;
 import java.util.List;
 import java.util.Map;
 import java.util.function.Supplier;
-import java.util.function.UnaryOperator;
 import java.util.stream.Collectors;
 
 import static 
org.apache.hudi.common.config.HoodieReaderConfig.RECORD_MERGE_IMPL_CLASSES_WRITE_CONFIG_KEY;
@@ -73,11 +68,6 @@ public class FlinkRowDataReaderContext extends 
HoodieReaderContext<RowData> {
   private final List<ExpressionPredicates.Predicate> predicates;
   private final Supplier<InternalSchemaManager> internalSchemaManager;
   private final HoodieTableConfig tableConfig;
-  // the converter is used to create a RowData contains primary key fields only
-  // for DELETE cases, it'll not be initialized if primary key semantics is 
lost.
-  // For e.g, if the pk fields are [a, b] but user only select a, then the pk
-  // semantics is lost.
-  private RecordKeyToRowDataConverter recordKeyRowConverter;
 
   public FlinkRowDataReaderContext(
       StorageConfiguration<?> storageConfiguration,
@@ -128,23 +118,6 @@ public class FlinkRowDataReaderContext extends 
HoodieReaderContext<RowData> {
     }
   }
 
-  @Override
-  public RowData seal(RowData rowData) {
-    if (rowData instanceof BinaryRowData) {
-      return ((BinaryRowData) rowData).copy();
-    }
-    return rowData;
-  }
-
-  @Override
-  public RowData toBinaryRow(Schema avroSchema, RowData record) {
-    if (record instanceof BinaryRowData) {
-      return record;
-    }
-    RowDataSerializer rowDataSerializer = 
RowDataAvroQueryContexts.getRowDataSerializer(avroSchema);
-    return rowDataSerializer.toBinaryRow(record);
-  }
-
   @Override
   public ClosableIterator<RowData> mergeBootstrapReaders(
       ClosableIterator<RowData> skeletonFileIterator,
@@ -191,25 +164,6 @@ public class FlinkRowDataReaderContext extends 
HoodieReaderContext<RowData> {
     };
   }
 
-  /**
-   * Creates a function that will reorder records of schema "from" to schema 
of "to".
-   * It's possible there exist fields in `to` schema, but not in `from` schema 
because of schema
-   * evolution.
-   *
-   * @param from           the schema of records to be passed into 
UnaryOperator
-   * @param to             the schema of records produced by UnaryOperator
-   * @param renamedColumns map of renamed columns where the key is the new 
name from the query and
-   *                       the value is the old name that exists in the file
-   * @return a function that takes in a record and returns the record with 
reordered columns
-   */
-  @Override
-  public UnaryOperator<RowData> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
-    RowType fromType = (RowType) 
RowDataAvroQueryContexts.fromAvroSchema(from).getRowType().getLogicalType();
-    RowType toType =  (RowType) 
RowDataAvroQueryContexts.fromAvroSchema(to).getRowType().getLogicalType();
-    RowProjection rowProjection = 
SchemaEvolvingRowDataProjection.instance(fromType, toType, renamedColumns);
-    return rowProjection::project;
-  }
-
   @Override
   public void setSchemaHandler(FileGroupReaderSchemaHandler<RowData> 
schemaHandler) {
     super.setSchemaHandler(schemaHandler);
@@ -229,7 +183,11 @@ public class FlinkRowDataReaderContext extends 
HoodieReaderContext<RowData> {
         .map(k -> 
Option.ofNullable(requiredSchema.getField(k)).map(Schema.Field::pos).orElse(-1))
         .mapToInt(Integer::intValue)
         .toArray();
-    recordKeyRowConverter = new RecordKeyToRowDataConverter(
+    // the converter is used to create a RowData contains primary key fields 
only
+    // for DELETE cases, it'll not be initialized if primary key semantics is 
lost.
+    // For e.g, if the pk fields are [a, b] but user only select a, then the pk
+    // semantics is lost.
+    RecordKeyToRowDataConverter recordKeyRowConverter = new 
RecordKeyToRowDataConverter(
         pkFieldsPos, (RowType) 
RowDataAvroQueryContexts.fromAvroSchema(requiredSchema).getRowType().getLogicalType());
     ((FlinkRecordContext) 
recordContext).setRecordKeyRowConverter(recordKeyRowConverter);
   }
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
index d77973111ff3..67d4e011a829 100644
--- 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
+++ 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
@@ -33,7 +33,6 @@ import 
org.apache.hudi.common.util.collection.ClosableIterator;
 import org.apache.hudi.common.util.collection.CloseableMappingIterator;
 import org.apache.hudi.common.util.collection.Pair;
 import org.apache.hudi.exception.HoodieAvroSchemaException;
-import org.apache.hudi.hadoop.utils.HoodieArrayWritableAvroUtils;
 import org.apache.hudi.io.storage.HoodieIOFactory;
 import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.StorageConfiguration;
@@ -60,14 +59,11 @@ import org.apache.hadoop.mapred.RecordReader;
 import org.apache.parquet.avro.AvroSchemaConverter;
 
 import java.io.IOException;
-import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Locale;
-import java.util.Map;
 import java.util.Set;
-import java.util.function.UnaryOperator;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
@@ -159,7 +155,7 @@ public class HiveHoodieReaderContext extends 
HoodieReaderContext<ArrayWritable>
       return recordIterator;
     }
     // record reader puts the required columns in the positions of the data 
schema and nulls the rest of the columns
-    return new CloseableMappingIterator<>(recordIterator, 
projectRecord(modifiedDataSchema, requiredSchema));
+    return new CloseableMappingIterator<>(recordIterator, 
recordContext.projectRecord(modifiedDataSchema, requiredSchema));
   }
 
   @Override
@@ -182,16 +178,6 @@ public class HiveHoodieReaderContext extends 
HoodieReaderContext<ArrayWritable>
     }
   }
 
-  @Override
-  public ArrayWritable seal(ArrayWritable record) {
-    return new ArrayWritable(Writable.class, Arrays.copyOf(record.get(), 
record.get().length));
-  }
-
-  @Override
-  public ArrayWritable toBinaryRow(Schema schema, ArrayWritable record) {
-    return record;
-  }
-
   @Override
   public ClosableIterator<ArrayWritable> 
mergeBootstrapReaders(ClosableIterator<ArrayWritable> skeletonFileIterator,
                                                                Schema 
skeletonRequiredSchema,
@@ -239,11 +225,6 @@ public class HiveHoodieReaderContext extends 
HoodieReaderContext<ArrayWritable>
     };
   }
 
-  @Override
-  public UnaryOperator<ArrayWritable> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
-    return record -> 
HoodieArrayWritableAvroUtils.rewriteRecordWithNewSchema(record, from, to, 
renamedColumns);
-  }
-
   public long getPos() throws IOException {
     if (firstRecordReader != null) {
       return firstRecordReader.getPos();
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveRecordContext.java 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveRecordContext.java
index 2a1d73b0b9c2..a2d64abba2e8 100644
--- a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveRecordContext.java
+++ b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveRecordContext.java
@@ -28,6 +28,7 @@ import org.apache.hudi.common.table.read.BufferedRecord;
 import org.apache.hudi.common.util.OrderingValues;
 import org.apache.hudi.hadoop.utils.HiveAvroSerializer;
 import org.apache.hudi.hadoop.utils.HiveJavaTypeConverter;
+import org.apache.hudi.hadoop.utils.HoodieArrayWritableAvroUtils;
 import org.apache.hudi.hadoop.utils.HoodieRealtimeRecordReaderUtils;
 
 import org.apache.avro.Schema;
@@ -43,8 +44,10 @@ import org.apache.hadoop.io.Text;
 import org.apache.hadoop.io.Writable;
 import org.apache.hadoop.io.WritableComparable;
 
+import java.util.Arrays;
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.function.UnaryOperator;
 
 public class HiveRecordContext extends RecordContext<ArrayWritable> {
 
@@ -132,4 +135,19 @@ public class HiveRecordContext extends 
RecordContext<ArrayWritable> {
   public ArrayWritable getDeleteRow(String recordKey) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
+
+  @Override
+  public ArrayWritable seal(ArrayWritable record) {
+    return new ArrayWritable(Writable.class, Arrays.copyOf(record.get(), 
record.get().length));
+  }
+
+  @Override
+  public ArrayWritable toBinaryRow(Schema schema, ArrayWritable record) {
+    return record;
+  }
+
+  @Override
+  public UnaryOperator<ArrayWritable> projectRecord(Schema from, Schema to, 
Map<String, String> renamedColumns) {
+    return record -> 
HoodieArrayWritableAvroUtils.rewriteRecordWithNewSchema(record, from, to, 
renamedColumns);
+  }
 }
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestBufferedRecordMerger.java
 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestBufferedRecordMerger.java
index ccc507aba344..5dbe1faca544 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestBufferedRecordMerger.java
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestBufferedRecordMerger.java
@@ -856,6 +856,22 @@ class TestBufferedRecordMerger extends 
SparkClientFunctionalTestHarness {
         throw new RuntimeException("Schema id is illegal: " + id);
       }
     }
+
+    @Override
+    public InternalRow seal(InternalRow record) {
+      return null;
+    }
+
+    @Override
+    public InternalRow toBinaryRow(Schema avroSchema, InternalRow record) {
+      return null;
+    }
+
+    @Override
+    public UnaryOperator<InternalRow> projectRecord(
+        Schema from, Schema to, Map<String, String> renamedColumns) {
+      return null;
+    }
   }
 
   static class DummyInternalRowReaderContext extends 
HoodieReaderContext<InternalRow> {
@@ -884,16 +900,6 @@ class TestBufferedRecordMerger extends 
SparkClientFunctionalTestHarness {
       return null;
     }
 
-    @Override
-    public InternalRow seal(InternalRow record) {
-      return null;
-    }
-
-    @Override
-    public InternalRow toBinaryRow(Schema avroSchema, InternalRow record) {
-      return null;
-    }
-
     @Override
     public ClosableIterator<InternalRow> mergeBootstrapReaders(
         ClosableIterator<InternalRow> skeletonFileIterator,
@@ -903,12 +909,6 @@ class TestBufferedRecordMerger extends 
SparkClientFunctionalTestHarness {
         List<Pair<String, Object>> requiredPartitionFieldAndValues) {
       return null;
     }
-
-    @Override
-    public UnaryOperator<InternalRow> projectRecord(
-        Schema from, Schema to, Map<String, String> renamedColumns) {
-      return null;
-    }
   }
 
   public static void assertRowEqual(InternalRow expected, InternalRow actual, 
Schema schema) {


Reply via email to