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 7dbd6f3c770 [HUDI-9322] Optimizing performance for hoodie log files 
writing (#13170)
7dbd6f3c770 is described below

commit 7dbd6f3c77004c72f8525efbb7662d1c3797c84b
Author: Shuo Cheng <[email protected]>
AuthorDate: Sun Apr 20 09:45:38 2025 +0800

    [HUDI-9322] Optimizing performance for hoodie log files writing (#13170)
    
    * [HUDI-9322] Optimizing performance for hoodie log files writing
    
    * eliminate the byte[] copy for avro data block serialization
    * add BaseAvroPayload#getRecordBytes
    * refactor the #bufferRecord of HoodieAppendHandle to make it more clear
    
    ---------
    
    Co-authored-by: danny0405 <[email protected]>
---
 .../org/apache/hudi/io/HoodieAppendHandle.java     |  89 +++++-----
 .../hudi/client/model/HoodieFlinkAvroRecord.java   | 187 ---------------------
 .../hudi/client/model/HoodieFlinkRecord.java       |  50 +++++-
 .../java/org/apache/hudi/io/FlinkAppendHandle.java |   5 +
 .../io/log/block/HoodieFlinkAvroDataBlock.java     |  52 ++++++
 .../io/log/block/HoodieFlinkParquetDataBlock.java  |   5 +-
 .../io/storage/row/RowDataParquetWriteSupport.java |   4 +-
 .../hudi/io/v2/FlinkRowDataHandleFactory.java      |  33 +---
 .../apache/hudi/io/v2/RowDataLogWriteHandle.java   |  30 +---
 .../apache/hudi/util/RowDataAvroQueryContexts.java | 113 +++++++++++++
 .../apache/hudi/util/RowDataToAvroConverters.java  |   4 +-
 .../java/org/apache/hudi/util/RowDataUtils.java    | 112 ++++++++++++
 .../hudi/common/model/HoodieSparkRecord.java       |  12 +-
 .../java/org/apache/hudi/avro/HoodieAvroUtils.java |   7 +
 .../hudi/common/config/HoodieStorageConfig.java    |   4 +-
 .../apache/hudi/common/model/BaseAvroPayload.java  |   4 +
 .../hudi/common/model/HoodieAvroIndexedRecord.java |  10 ++
 .../apache/hudi/common/model/HoodieAvroRecord.java |  15 ++
 .../hudi/common/model/HoodieEmptyRecord.java       |  10 ++
 .../org/apache/hudi/common/model/HoodieRecord.java |   6 +
 .../model/HoodieRecordCompatibilityInterface.java  |   2 +
 .../table/log/block/HoodieAvroDataBlock.java       |  28 ++-
 .../common/table/log/block/HoodieCommandBlock.java |   5 +-
 .../common/table/log/block/HoodieCorruptBlock.java |   5 +-
 .../common/table/log/block/HoodieDataBlock.java    |  10 +-
 .../common/table/log/block/HoodieDeleteBlock.java  |  11 +-
 .../table/log/block/HoodieHFileDataBlock.java      |   3 +-
 .../common/table/log/block/HoodieLogBlock.java     |  17 +-
 .../table/log/block/HoodieParquetDataBlock.java    |   3 +-
 .../apache/hudi/common/util/FileFormatUtils.java   |  25 +--
 .../hudi/metadata/HoodieTableMetadataUtil.java     |  15 +-
 .../table/log/block/TestHoodieAvroDataBlock.java   |   9 +-
 .../hudi/sink/RowDataStreamWriteFunction.java      |   2 +-
 .../hudi/sink/transform/RecordConverter.java       |  57 +------
 .../apache/hudi/util/OrderingValueExtractor.java   |  57 -------
 .../apache/hudi/table/ITTestHoodieDataSource.java  |  13 --
 .../common/table/log/HoodieLogFormatWriter.java    |   9 +-
 .../org/apache/hudi/common/util/HFileUtils.java    |   6 +-
 .../java/org/apache/hudi/common/util/OrcUtils.java |  27 +--
 .../org/apache/hudi/common/util/ParquetUtils.java  |  22 +--
 .../common/functional/TestHoodieLogFormat.java     |   2 +-
 .../table/log/block/TestHoodieDeleteBlock.java     |   2 +-
 .../org/apache/hudi/hadoop/HoodieHiveRecord.java   |  10 ++
 .../hudi/testutils/LogFileColStatsTestUtil.java    |   2 +-
 44 files changed, 598 insertions(+), 496 deletions(-)

diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
index 7571d173aaa..c35b4aa98c3 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
@@ -284,7 +284,7 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
     return hoodieRecord.getCurrentLocation() != null;
   }
 
-  private Option<HoodieRecord> prepareRecord(HoodieRecord<T> hoodieRecord) {
+  private void bufferRecord(HoodieRecord<T> hoodieRecord) {
     Option<Map<String, String>> recordMetadata = hoodieRecord.getMetadata();
     Schema schema = useWriterSchema ? writeSchemaWithMetaFields : writeSchema;
     try {
@@ -293,36 +293,11 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
       boolean isUpdateRecord = isUpdateRecord(hoodieRecord);
       
recordProperties.put(HoodiePayloadProps.PAYLOAD_IS_UPDATE_RECORD_FOR_MOR, 
String.valueOf(isUpdateRecord));
 
-      final Option<HoodieRecord> finalRecordOpt;
       // Check for delete
       if (!hoodieRecord.isDelete(schema, recordProperties) || 
config.allowOperationMetadataField()) {
-        // Check if the record should be ignored (special case for 
[[ExpressionPayload]])
-        if (hoodieRecord.shouldIgnore(schema, recordProperties)) {
-          return Option.of(hoodieRecord);
-        }
-
-        // Prepend meta-fields into the record
-        MetadataValues metadataValues = populateMetadataFields(hoodieRecord);
-        HoodieRecord populatedRecord =
-            hoodieRecord.prependMetaFields(schema, writeSchemaWithMetaFields, 
metadataValues, recordProperties);
-
-        // NOTE: Record have to be cloned here to make sure if it holds 
low-level engine-specific
-        //       payload pointing into a shared, mutable (underlying) buffer 
we get a clean copy of
-        //       it since these records will be put into the recordList(List).
-        finalRecordOpt = Option.of(populatedRecord.copy());
-        if (isUpdateRecord || isLogCompaction) {
-          updatedRecordsWritten++;
-        } else {
-          insertRecordsWritten++;
-        }
-        recordsWritten++;
+        bufferInsertAndUpdate(schema, hoodieRecord, isUpdateRecord);
       } else {
-        finalRecordOpt = Option.empty();
-        // Clear the new location as the record was deleted
-        hoodieRecord.unseal();
-        hoodieRecord.clearNewLocation();
-        hoodieRecord.seal();
-        recordsDeleted++;
+        bufferDelete(hoodieRecord);
       }
 
       writeStatus.markSuccess(hoodieRecord, recordMetadata);
@@ -330,12 +305,10 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
       // part of marking
       // record successful.
       hoodieRecord.deflate();
-      return finalRecordOpt;
     } catch (Exception e) {
-      LOG.error("Error writing record  " + hoodieRecord, e);
+      LOG.error("Error writing record {}", hoodieRecord, e);
       writeStatus.markFailure(hoodieRecord, e, recordMetadata);
     }
-    return Option.empty();
   }
 
   private MetadataValues populateMetadataFields(HoodieRecord<T> hoodieRecord) {
@@ -459,7 +432,7 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
           .map(fieldName -> 
HoodieAvroUtils.getSchemaForField(writeSchemaWithMetaFields, 
fieldName)).collect(Collectors.toList());
       try {
         Map<String, HoodieColumnRangeMetadata<Comparable>> 
columnRangeMetadataMap =
-            collectColumnRangeMetadata(recordList, fieldsToIndex, 
stat.getPath(), writeSchemaWithMetaFields);
+            collectColumnRangeMetadata(recordList, fieldsToIndex, 
stat.getPath(), writeSchemaWithMetaFields, storage.getConf());
         stat.putRecordsStats(columnRangeMetadataMap);
       } catch (HoodieException e) {
         throw new HoodieAppendException("Failed to extract append result", e);
@@ -625,29 +598,49 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
       record.seal();
     }
     // fetch the ordering val first in case the record was deflated.
-    final Comparable<?> orderingVal = record.getOrderingValue(writeSchema, 
recordProperties);
-    Option<HoodieRecord> indexedRecord = prepareRecord(record);
-    if (indexedRecord.isPresent()) {
-      // Skip the ignored record.
-      try {
-        if (!indexedRecord.get().shouldIgnore(writeSchema, recordProperties)) {
-          recordList.add(indexedRecord.get());
-        }
-      } catch (IOException e) {
-        writeStatus.markFailure(record, e, record.getMetadata());
-        LOG.error("Error writing record  " + indexedRecord.get(), e);
-      }
+    bufferRecord(record);
+    numberOfRecords++;
+  }
+
+  private void bufferInsertAndUpdate(Schema schema, HoodieRecord<T> 
hoodieRecord, boolean isUpdateRecord) throws IOException {
+    // Check if the record should be ignored (special case for 
[[ExpressionPayload]])
+    if (hoodieRecord.shouldIgnore(schema, recordProperties)) {
+      return;
+    }
+
+    // Prepend meta-fields into the record
+    MetadataValues metadataValues = populateMetadataFields(hoodieRecord);
+    HoodieRecord populatedRecord =
+        hoodieRecord.prependMetaFields(schema, writeSchemaWithMetaFields, 
metadataValues, recordProperties);
+
+    // NOTE: Record have to be cloned here to make sure if it holds low-level 
engine-specific
+    //       payload pointing into a shared, mutable (underlying) buffer we 
get a clean copy of
+    //       it since these records will be put into the recordList(List).
+    recordList.add(populatedRecord.copy());
+    if (isUpdateRecord || isLogCompaction) {
+      updatedRecordsWritten++;
     } else {
-      long position = baseFileInstantTimeOfPositions.isPresent() ? 
record.getCurrentPosition() : -1L;
-      
recordsToDeleteWithPositions.add(Pair.of(DeleteRecord.create(record.getKey(), 
orderingVal), position));
+      insertRecordsWritten++;
     }
-    numberOfRecords++;
+    recordsWritten++;
+  }
+
+  private void bufferDelete(HoodieRecord<T> hoodieRecord) {
+    // Clear the new location as the record was deleted
+    hoodieRecord.unseal();
+    hoodieRecord.clearNewLocation();
+    hoodieRecord.seal();
+    recordsDeleted++;
+
+    final Comparable<?> orderingVal = 
hoodieRecord.getOrderingValue(writeSchema, recordProperties);
+    long position = baseFileInstantTimeOfPositions.isPresent() ? 
hoodieRecord.getCurrentPosition() : -1L;
+    
recordsToDeleteWithPositions.add(Pair.of(DeleteRecord.create(hoodieRecord.getKey(),
 orderingVal), position));
   }
 
   /**
    * Checks if the number of records have reached the set threshold and then 
flushes the records to disk.
    */
-  private void flushToDiskIfRequired(HoodieRecord record, boolean 
appendDeleteBlocks) {
+  protected void flushToDiskIfRequired(HoodieRecord record, boolean 
appendDeleteBlocks) {
     if (numberOfRecords >= (int) (maxBlockSize / averageRecordSize)
         || numberOfRecords % NUMBER_OF_RECORDS_TO_ESTIMATE_RECORD_SIZE == 0) {
       averageRecordSize = (long) (averageRecordSize * 0.8 + 
sizeEstimator.sizeEstimate(record) * 0.2);
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkAvroRecord.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkAvroRecord.java
deleted file mode 100644
index a8c9204775d..00000000000
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkAvroRecord.java
+++ /dev/null
@@ -1,187 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.hudi.client.model;
-
-import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
-import org.apache.hudi.common.model.HoodieKey;
-import org.apache.hudi.common.model.HoodieOperation;
-import org.apache.hudi.common.model.HoodieRecord;
-import org.apache.hudi.common.model.MetadataValues;
-import org.apache.hudi.common.util.Option;
-import org.apache.hudi.common.util.collection.Pair;
-import org.apache.hudi.keygen.BaseKeyGenerator;
-
-import com.esotericsoftware.kryo.Kryo;
-import com.esotericsoftware.kryo.io.Input;
-import com.esotericsoftware.kryo.io.Output;
-import org.apache.avro.Schema;
-import org.apache.avro.generic.GenericData;
-import org.apache.avro.generic.GenericRecord;
-import org.apache.avro.generic.IndexedRecord;
-
-import java.util.Map;
-import java.util.Properties;
-
-/**
- * Flink implementation of `HoodieRecord`, which is expected to hold Avro 
{@code IndexedRecord} as payload.
- * It's only used by writer when the log block format type is AVRO.
- */
-public class HoodieFlinkAvroRecord extends HoodieRecord<IndexedRecord> {
-  private Comparable<?> orderingValue = 0;
-
-  public HoodieFlinkAvroRecord(HoodieKey key, HoodieOperation op, 
Comparable<?> orderingValue, IndexedRecord record) {
-    super(key, record, op, Option.empty());
-    this.orderingValue = orderingValue;
-  }
-
-  @Override
-  public HoodieRecord<IndexedRecord> newInstance() {
-    return new HoodieFlinkAvroRecord(key, operation, orderingValue, data);
-  }
-
-  @Override
-  public HoodieRecord<IndexedRecord> newInstance(HoodieKey key, 
HoodieOperation op) {
-    return new HoodieFlinkAvroRecord(key, op, orderingValue, data);
-  }
-
-  @Override
-  public HoodieRecord<IndexedRecord> newInstance(HoodieKey key) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public Comparable<?> getOrderingValue(Schema recordSchema, Properties props) 
{
-    return this.orderingValue;
-  }
-
-  @Override
-  public HoodieRecordType getRecordType() {
-    return HoodieRecordType.AVRO;
-  }
-
-  @Override
-  public String getRecordKey(Schema recordSchema, Option<BaseKeyGenerator> 
keyGeneratorOpt) {
-    return getRecordKey();
-  }
-
-  @Override
-  public String getRecordKey(Schema recordSchema, String keyFieldName) {
-    return getRecordKey();
-  }
-
-  @Override
-  protected void writeRecordPayload(IndexedRecord payload, Kryo kryo, Output 
output) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  protected IndexedRecord readRecordPayload(Kryo kryo, Input input) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public Object[] getColumnValues(Schema recordSchema, String[] columns, 
boolean consistentLogicalTimestampEnabled) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public HoodieRecord joinWith(HoodieRecord other, Schema targetSchema) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public HoodieRecord prependMetaFields(Schema recordSchema, Schema 
targetSchema, MetadataValues metadataValues, Properties props) {
-    GenericRecord recordWithMetaFields = new GenericData.Record(targetSchema);
-    // update meta fields
-    if (!metadataValues.isEmpty()) {
-      String[] values = metadataValues.getValues();
-      for (int i = 0; i < values.length; i++) {
-        if (values[i] != null) {
-          recordWithMetaFields.put(i, values[i]);
-        }
-      }
-    }
-    // update data fields
-    int metaFieldsSize = targetSchema.getFields().size() - 
recordSchema.getFields().size();
-    for (int i = 0; i < recordSchema.getFields().size(); i++) {
-      recordWithMetaFields.put(metaFieldsSize + i, data.get(i));
-    }
-    return new HoodieFlinkAvroRecord(key, operation, orderingValue, 
recordWithMetaFields);
-  }
-
-  @Override
-  public HoodieRecord rewriteRecordWithNewSchema(Schema recordSchema, 
Properties props, Schema newSchema, Map<String, String> renameCols) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public boolean isDelete(Schema recordSchema, Properties props) {
-    if (data == null) {
-      return true;
-    }
-
-    if (HoodieOperation.isDelete(getOperation())) {
-      return true;
-    }
-
-    // Use data field to decide.
-    Schema.Field deleteField = recordSchema.getField(HOODIE_IS_DELETED_FIELD);
-    if (deleteField == null) {
-      return false;
-    }
-    Object deleteMarker = data.get(deleteField.pos());
-    return deleteMarker instanceof Boolean && (Boolean) deleteMarker;
-  }
-
-  @Override
-  public boolean shouldIgnore(Schema recordSchema, Properties props) {
-    return false;
-  }
-
-  @Override
-  public HoodieRecord<IndexedRecord> copy() {
-    return this;
-  }
-
-  @Override
-  public Option<Map<String, String>> getMetadata() {
-    return Option.empty();
-  }
-
-  @Override
-  public HoodieRecord wrapIntoHoodieRecordPayloadWithParams(Schema 
recordSchema, Properties props, Option<Pair<String, String>> 
simpleKeyGenFieldsOpt, Boolean withOperation,
-                                                            Option<String> 
partitionNameOp, Boolean populateMetaFieldsOp, Option<Schema> 
schemaWithoutMetaFields) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public HoodieRecord wrapIntoHoodieRecordPayloadWithKeyGen(Schema 
recordSchema, Properties props, Option<BaseKeyGenerator> keyGen) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public HoodieRecord truncateRecordKey(Schema recordSchema, Properties props, 
String keyFieldName) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
-  }
-
-  @Override
-  public Option<HoodieAvroIndexedRecord> toIndexedRecord(Schema recordSchema, 
Properties props) {
-    return Option.of(new HoodieAvroIndexedRecord(getKey(), getData(), 
getOperation(), getMetadata()));
-  }
-}
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkRecord.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkRecord.java
index 77f0456b17f..7f56f0b0cb2 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkRecord.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkRecord.java
@@ -18,38 +18,49 @@
 
 package org.apache.hudi.client.model;
 
+import org.apache.hudi.avro.HoodieAvroUtils;
+import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
 import org.apache.hudi.common.model.HoodieKey;
 import org.apache.hudi.common.model.HoodieOperation;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.MetadataValues;
+import org.apache.hudi.common.util.ConfigUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.Pair;
 import org.apache.hudi.keygen.BaseKeyGenerator;
+import org.apache.hudi.util.RowDataAvroQueryContexts;
+import org.apache.hudi.util.RowDataAvroQueryContexts.RowDataQueryContext;
 
 import com.esotericsoftware.kryo.Kryo;
 import com.esotericsoftware.kryo.io.Input;
 import com.esotericsoftware.kryo.io.Output;
 import org.apache.avro.Schema;
+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.StringData;
 import org.apache.flink.table.data.utils.JoinedRowData;
 
-import java.io.IOException;
 import java.util.Map;
 import java.util.Properties;
 
+import static org.apache.hudi.common.util.StringUtils.isNullOrEmpty;
+
 /**
  * Flink Engine-specific Implementations of `HoodieRecord`, which is expected 
to hold {@code RowData} as payload.
  */
 public class HoodieFlinkRecord extends HoodieRecord<RowData> {
-  private Comparable<?> orderingValue = 0;
+  private Comparable<?> orderingValue;
 
   public HoodieFlinkRecord(RowData rowData) {
     super(null, rowData);
   }
 
+  public HoodieFlinkRecord(HoodieKey key, HoodieOperation op, RowData rowData) 
{
+    super(key, rowData, op, Option.empty());
+  }
+
   public HoodieFlinkRecord(HoodieKey key, HoodieOperation op, Comparable<?> 
orderingValue, RowData rowData) {
     super(key, rowData, op, Option.empty());
     this.orderingValue = orderingValue;
@@ -72,6 +83,16 @@ public class HoodieFlinkRecord extends HoodieRecord<RowData> 
{
 
   @Override
   public Comparable<?> getOrderingValue(Schema recordSchema, Properties props) 
{
+    if (this.orderingValue == null) {
+      String orderingField = ConfigUtils.getOrderingField(props);
+      if (isNullOrEmpty(orderingField)) {
+        this.orderingValue = DEFAULT_ORDERING_VALUE;
+      } else {
+        boolean utcTimezone = 
Boolean.parseBoolean(props.getProperty("read.utc-timezone", "true"));
+        RowDataAvroQueryContexts.FieldQueryContext context = 
RowDataAvroQueryContexts.fromAvroSchema(recordSchema, 
utcTimezone).getFieldQueryContext(orderingField);
+        this.orderingValue = (Comparable<?>) context.getValAsJava(this.data);
+      }
+    }
     return this.orderingValue;
   }
 
@@ -105,6 +126,14 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
 
+  @Override
+  public Object getColumnValueAsJava(Schema recordSchema, String columns, 
Properties props) {
+    boolean utcTimezone = Boolean.parseBoolean(props.getProperty(
+        HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(), 
HoodieStorageConfig.WRITE_UTC_TIMEZONE.defaultValue().toString()));
+    RowDataQueryContext rowDataQueryContext = 
RowDataAvroQueryContexts.fromAvroSchema(recordSchema, utcTimezone);
+    return 
rowDataQueryContext.getFieldQueryContext(columns).getValAsJava(data);
+  }
+
   @Override
   public HoodieRecord joinWith(HoodieRecord other, Schema targetSchema) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
@@ -144,7 +173,7 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
   }
 
   @Override
-  public boolean shouldIgnore(Schema recordSchema, Properties props) throws 
IOException {
+  public boolean shouldIgnore(Schema recordSchema, Properties props) {
     return false;
   }
 
@@ -176,6 +205,19 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
 
   @Override
   public Option<HoodieAvroIndexedRecord> toIndexedRecord(Schema recordSchema, 
Properties props) {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
+    boolean utcTimezone = Boolean.parseBoolean(props.getProperty(
+        HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(), 
HoodieStorageConfig.WRITE_UTC_TIMEZONE.defaultValue().toString()));
+    RowDataQueryContext rowDataQueryContext = 
RowDataAvroQueryContexts.fromAvroSchema(recordSchema, utcTimezone);
+    IndexedRecord indexedRecord = (IndexedRecord) 
rowDataQueryContext.getRowDataToAvroConverter().convert(recordSchema, 
getData());
+    return Option.of(new HoodieAvroIndexedRecord(getKey(), indexedRecord, 
getOperation(), getMetadata()));
+  }
+
+  @Override
+  public byte[] getAvroBytes(Schema recordSchema, Properties props) {
+    boolean utcTimezone = Boolean.parseBoolean(props.getProperty(
+        HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(), 
HoodieStorageConfig.WRITE_UTC_TIMEZONE.defaultValue().toString()));
+    RowDataQueryContext rowDataQueryContext = 
RowDataAvroQueryContexts.fromAvroSchema(recordSchema, utcTimezone);
+    IndexedRecord indexedRecord = (IndexedRecord) 
rowDataQueryContext.getRowDataToAvroConverter().convert(recordSchema, 
getData());
+    return HoodieAvroUtils.avroToBytes(indexedRecord);
   }
 }
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkAppendHandle.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkAppendHandle.java
index 35b988c142b..2564d94bd3a 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkAppendHandle.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkAppendHandle.java
@@ -67,6 +67,11 @@ public class FlinkAppendHandle<T, I, K, O>
     this.bucketType = bucketType;
   }
 
+  @Override
+  protected void flushToDiskIfRequired(HoodieRecord record, boolean 
appendDeleteBlocks) {
+    // do not flush for one batch of records
+  }
+
   @Override
   protected void createMarkerFile(String partitionPath, String dataFileName) {
     // In some rare cases, the task was pulled up again with same write file 
name,
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/log/block/HoodieFlinkAvroDataBlock.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/log/block/HoodieFlinkAvroDataBlock.java
new file mode 100644
index 00000000000..0f7a0a76166
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/log/block/HoodieFlinkAvroDataBlock.java
@@ -0,0 +1,52 @@
+/*
+ * 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.log.block;
+
+import org.apache.hudi.common.config.HoodieStorageConfig;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.table.log.block.HoodieAvroDataBlock;
+import org.apache.hudi.storage.StorageConfiguration;
+
+import org.jetbrains.annotations.NotNull;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+
+/**
+ * HoodieFlinkAvroDataBlock contains a list of records serialized using Avro. 
It is used with the Parquet base file format.
+ */
+public class HoodieFlinkAvroDataBlock extends HoodieAvroDataBlock {
+
+  public HoodieFlinkAvroDataBlock(
+      @NotNull List<HoodieRecord> records,
+      @NotNull Map<HeaderMetadataType, String> header,
+      @NotNull String keyField) {
+    super(records, header, keyField);
+  }
+
+  @Override
+  protected Properties initProperties(StorageConfiguration<?> storageConfig) {
+    Properties properties = new Properties();
+    properties.setProperty(
+        HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(),
+        storageConfig.getString(HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(), 
HoodieStorageConfig.WRITE_UTC_TIMEZONE.defaultValue().toString()));
+    return properties;
+  }
+}
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/log/block/HoodieFlinkParquetDataBlock.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/log/block/HoodieFlinkParquetDataBlock.java
index f8e69eeff94..bebe3f807af 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/log/block/HoodieFlinkParquetDataBlock.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/log/block/HoodieFlinkParquetDataBlock.java
@@ -34,6 +34,7 @@ import org.apache.hudi.storage.HoodieStorage;
 import org.apache.avro.Schema;
 import org.apache.parquet.hadoop.metadata.ParquetMetadata;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.HashMap;
 import java.util.Iterator;
@@ -71,7 +72,7 @@ public class HoodieFlinkParquetDataBlock extends 
HoodieParquetDataBlock implemen
   }
 
   @Override
-  public byte[] getContentBytes(HoodieStorage storage) throws IOException {
+  public ByteArrayOutputStream getContentBytes(HoodieStorage storage) throws 
IOException {
     Map<String, String> paramsMap = new HashMap<>();
     paramsMap.put(PARQUET_COMPRESSION_CODEC_NAME.key(), 
compressionCodecName.get());
     paramsMap.put(PARQUET_COMPRESSION_RATIO_FRACTION.key(), 
String.valueOf(expectedCompressionRatio.get()));
@@ -79,7 +80,7 @@ public class HoodieFlinkParquetDataBlock extends 
HoodieParquetDataBlock implemen
     Schema writerSchema = AvroSchemaCache.intern(new Schema.Parser().parse(
         
super.getLogBlockHeader().get(HoodieLogBlock.HeaderMetadataType.SCHEMA)));
 
-    Pair<byte[], Object> result =
+    Pair<ByteArrayOutputStream, Object> result =
         HoodieIOFactory.getIOFactory(storage).getFileFormatUtils(PARQUET)
             .serializeRecordsToLogBlock(
                 storage,
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
index f2c87cb25ff..7315461db50 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
@@ -58,8 +58,8 @@ public class RowDataParquetWriteSupport extends 
WriteSupport<RowData> {
     // should make the utc timestamp configurable
     boolean utcTimestamp =
         hadoopConf.getBoolean(
-            HoodieStorageConfig.PARQUET_WRITE_UTC_TIMEZONE.key(),
-            HoodieStorageConfig.PARQUET_WRITE_UTC_TIMEZONE.defaultValue());
+            HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(),
+            HoodieStorageConfig.WRITE_UTC_TIMEZONE.defaultValue());
     this.writer = new ParquetRowDataWriter(recordConsumer, rowType, schema, 
utcTimestamp);
   }
 
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/FlinkRowDataHandleFactory.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/FlinkRowDataHandleFactory.java
index 27dfd6f6fd1..22788cda57c 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/FlinkRowDataHandleFactory.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/FlinkRowDataHandleFactory.java
@@ -21,14 +21,11 @@ package org.apache.hudi.io.v2;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieTableType;
 import org.apache.hudi.common.table.HoodieTableConfig;
-import 
org.apache.hudi.common.table.log.block.HoodieLogBlock.HoodieLogBlockType;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieException;
-import org.apache.hudi.io.FlinkAppendHandle;
 import org.apache.hudi.io.HoodieWriteHandle;
 import org.apache.hudi.table.HoodieTable;
 import org.apache.hudi.table.action.commit.BucketInfo;
-import org.apache.hudi.util.CommonClientUtils;
 
 import org.apache.hadoop.fs.Path;
 
@@ -69,27 +66,15 @@ public class FlinkRowDataHandleFactory {
         String instantTime,
         HoodieTable<T, I, K, O> table,
         Iterator<HoodieRecord<T>> recordIterator) {
-      if (CommonClientUtils.getLogBlockType(config, 
table.getMetaClient().getTableConfig()) == 
HoodieLogBlockType.PARQUET_DATA_BLOCK) {
-        return new RowDataLogWriteHandle<>(
-            config,
-            instantTime,
-            table,
-            recordIterator,
-            bucketInfo.getFileIdPrefix(),
-            bucketInfo.getPartitionPath(),
-            bucketInfo.getBucketType(),
-            table.getTaskContextSupplier());
-      } else {
-        return new FlinkAppendHandle<>(
-            config,
-            instantTime,
-            table,
-            bucketInfo.getPartitionPath(),
-            bucketInfo.getFileIdPrefix(),
-            bucketInfo.getBucketType(),
-            recordIterator,
-            table.getTaskContextSupplier());
-      }
+      return new RowDataLogWriteHandle<>(
+          config,
+          instantTime,
+          table,
+          recordIterator,
+          bucketInfo.getFileIdPrefix(),
+          bucketInfo.getPartitionPath(),
+          bucketInfo.getBucketType(),
+          table.getTaskContextSupplier());
     }
   }
 
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/RowDataLogWriteHandle.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/RowDataLogWriteHandle.java
index 2682d718f30..71b9ed3624a 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/RowDataLogWriteHandle.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/v2/RowDataLogWriteHandle.java
@@ -34,6 +34,7 @@ import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.io.FlinkAppendHandle;
 import org.apache.hudi.io.MiniBatchHandle;
+import org.apache.hudi.io.log.block.HoodieFlinkAvroDataBlock;
 import org.apache.hudi.io.log.block.HoodieFlinkParquetDataBlock;
 import org.apache.hudi.io.storage.ColumnRangeMetadataProvider;
 import org.apache.hudi.io.storage.row.HoodieFlinkIOFactory;
@@ -88,8 +89,8 @@ public class RowDataLogWriteHandle<T, I, K, O>
 
   private void initWriteConf(StorageConfiguration<?> storageConf, 
HoodieWriteConfig writeConfig) {
     storageConf.set(
-        HoodieStorageConfig.PARQUET_WRITE_UTC_TIMEZONE.key(),
-        
writeConfig.getString(HoodieStorageConfig.PARQUET_WRITE_UTC_TIMEZONE.key()));
+        HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(),
+        writeConfig.getString(HoodieStorageConfig.WRITE_UTC_TIMEZONE.key()));
     storageConf.set(
         HoodieStorageConfig.HOODIE_IO_FACTORY_CLASS.key(),
         HoodieFlinkIOFactory.class.getName());
@@ -100,25 +101,6 @@ public class RowDataLogWriteHandle<T, I, K, O>
     return new FlinkRecordSizeEstimator();
   }
 
-  @Override
-  protected HoodieLogBlockType getLogBlockType() {
-    Option<HoodieLogBlock.HoodieLogBlockType> logBlockTypeOpt = 
config.getLogDataBlockFormat();
-    if (logBlockTypeOpt.isPresent()) {
-      return logBlockTypeOpt.get();
-    }
-    // Fallback to deduce data-block type based on the base file format
-    switch (hoodieTable.getBaseFileFormat()) {
-      case PARQUET:
-      case ORC:
-        return HoodieLogBlockType.PARQUET_DATA_BLOCK;
-      case HFILE:
-        return HoodieLogBlock.HoodieLogBlockType.HFILE_DATA_BLOCK;
-      default:
-        throw new HoodieException("Base file format " + 
hoodieTable.getBaseFileFormat()
-            + " does not have associated log block type");
-    }
-  }
-
   /**
    * Flink writer does not support record-position for update/delete 
currently, will be supported later, see HUDI-9192.
    */
@@ -129,6 +111,10 @@ public class RowDataLogWriteHandle<T, I, K, O>
 
   @Override
   protected void processAppendResult(AppendResult result, 
Option<HoodieLogBlock> dataBlock) {
+    if (getLogBlockType() == HoodieLogBlockType.AVRO_DATA_BLOCK) {
+      super.processAppendResult(result, dataBlock);
+      return;
+    }
     HoodieDeltaWriteStat stat = (HoodieDeltaWriteStat) 
this.writeStatus.getStat();
     updateWriteStatus(result, stat);
 
@@ -194,6 +180,8 @@ public class RowDataLogWriteHandle<T, I, K, O>
             writeConfig.getParquetCompressionCodec(),
             writeConfig.getParquetCompressionRatio(),
             writeConfig.parquetDictionaryEnabled());
+      case AVRO_DATA_BLOCK:
+        return new HoodieFlinkAvroDataBlock(records, header, keyField);
       default:
         throw new HoodieException("Data block format " + logDataBlockFormat + 
" is not implemented for Flink RowData append handle.");
     }
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataAvroQueryContexts.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataAvroQueryContexts.java
new file mode 100644
index 00000000000..2204f7dcdf2
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataAvroQueryContexts.java
@@ -0,0 +1,113 @@
+/*
+ * 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.util;
+
+import org.apache.avro.Schema;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.RowType;
+import org.apache.hudi.util.RowDataToAvroConverters.RowDataToAvroConverter;
+
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.function.Function;
+
+/**
+ * Maintains auxiliary utilities for row data fields handling.
+ */
+public class RowDataAvroQueryContexts {
+  private static final Map<Schema, RowDataQueryContext> QUERY_CONTEXT_MAP = 
new ConcurrentHashMap<>();
+
+  public static RowDataQueryContext fromAvroSchema(Schema avroSchema, boolean 
utcTimezone) {
+    return QUERY_CONTEXT_MAP.computeIfAbsent(avroSchema, k -> {
+      DataType dataType = AvroSchemaConverter.convertToDataType(avroSchema);
+      RowType rowType = (RowType) dataType.getLogicalType();
+      RowType.RowField[] rowFields = rowType.getFields().toArray(new 
RowType.RowField[0]);
+      RowData.FieldGetter[] fieldGetters = new 
RowData.FieldGetter[rowFields.length];
+      Map<String, FieldQueryContext> contextMap = new HashMap<>();
+      for (int i = 0; i < rowFields.length; i++) {
+        LogicalType fieldType = rowFields[i].getType();
+        RowData.FieldGetter fieldGetter = 
RowData.createFieldGetter(rowFields[i].getType(), i);
+        fieldGetters[i] = fieldGetter;
+        contextMap.put(rowFields[i].getName(), 
FieldQueryContext.create(fieldType, fieldGetter, utcTimezone));
+      }
+      RowDataToAvroConverter rowDataToAvroConverter = 
RowDataToAvroConverters.createConverter(rowType, utcTimezone);
+      return RowDataQueryContext.create(contextMap, fieldGetters, 
rowDataToAvroConverter);
+    });
+  }
+
+  public static class RowDataQueryContext {
+    private final Map<String, FieldQueryContext> contextMap;
+    private final RowData.FieldGetter[] fieldGetters;
+    private final RowDataToAvroConverter rowDataToAvroConverter;
+    private RowDataQueryContext(Map<String, FieldQueryContext> contextMap, 
RowData.FieldGetter[] fieldGetters, RowDataToAvroConverter 
rowDataAvroConverter) {
+      this.contextMap = contextMap;
+      this.fieldGetters = fieldGetters;
+      this.rowDataToAvroConverter = rowDataAvroConverter;
+    }
+
+    public static RowDataQueryContext create(
+        Map<String, FieldQueryContext> contextMap,
+        RowData.FieldGetter[] fieldGetters,
+        RowDataToAvroConverter rowDataToAvroConverter) {
+      return new RowDataQueryContext(contextMap, fieldGetters, 
rowDataToAvroConverter);
+    }
+
+    public FieldQueryContext getFieldQueryContext(String fieldName) {
+      return contextMap.get(fieldName);
+    }
+
+    public RowData.FieldGetter[] fieldGetters() {
+      return fieldGetters;
+    }
+
+    public RowDataToAvroConverter getRowDataToAvroConverter() {
+      return rowDataToAvroConverter;
+    }
+  }
+
+  public static class FieldQueryContext {
+    private final LogicalType logicalType;
+    private final RowData.FieldGetter fieldGetter;
+    private final Function<Object, Object> javaTypeConverter;
+    private FieldQueryContext(LogicalType logicalType, RowData.FieldGetter 
fieldGetter, boolean utcTimezone) {
+      this.logicalType = logicalType;
+      this.fieldGetter = fieldGetter;
+      this.javaTypeConverter = RowDataUtils.orderingValFunc(logicalType, 
utcTimezone);
+    }
+
+    public static FieldQueryContext create(LogicalType logicalType, 
RowData.FieldGetter fieldGetter, boolean utcTimezone) {
+      return new FieldQueryContext(logicalType, fieldGetter, utcTimezone);
+    }
+
+    public LogicalType getLogicalType() {
+      return logicalType;
+    }
+
+    public RowData.FieldGetter getFieldGetter() {
+      return fieldGetter;
+    }
+
+    public Object getValAsJava(RowData rowData) {
+      return this.javaTypeConverter.apply(fieldGetter.getFieldOrNull(rowData));
+    }
+  }
+}
\ No newline at end of file
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/RowDataToAvroConverters.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataToAvroConverters.java
similarity index 99%
rename from 
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/RowDataToAvroConverters.java
rename to 
hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataToAvroConverters.java
index f54abd4a16b..4c9216fef68 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/RowDataToAvroConverters.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataToAvroConverters.java
@@ -160,7 +160,7 @@ public class RowDataToAvroConverters {
             };
         break;
       case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
-        int precision = DataTypeUtils.precision(type);
+        int precision = RowDataUtils.precision(type);
         if (precision <= 3) {
           converter = new RowDataToAvroConverter() {
             private static final long serialVersionUID = 1L;
@@ -185,7 +185,7 @@ public class RowDataToAvroConverters {
         }
         break;
       case TIMESTAMP_WITHOUT_TIME_ZONE:
-        precision = DataTypeUtils.precision(type);
+        precision = RowDataUtils.precision(type);
         if (precision <= 3) {
           converter =
               new RowDataToAvroConverter() {
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataUtils.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataUtils.java
new file mode 100644
index 00000000000..17716ca70ab
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/RowDataUtils.java
@@ -0,0 +1,112 @@
+/*
+ * 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.util;
+
+import org.apache.hudi.exception.HoodieValidationException;
+
+import org.apache.flink.table.data.DecimalData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.data.TimestampData;
+import org.apache.flink.table.types.logical.LocalZonedTimestampType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.TimestampType;
+
+import java.nio.ByteBuffer;
+import java.time.Instant;
+import java.time.LocalDate;
+import java.util.function.Function;
+
+/**
+ * Utils for get/set operations on {@link RowData}.
+ */
+public class RowDataUtils {
+  /**
+   * Resolve the native Java object from given row data field value.
+   *
+   * <p>IMPORTANT: the logic references the row-data to avro conversion in 
{@code RowDataToAvroConverters.createConverter}
+   * and {@code HoodieAvroUtils.convertValueForAvroLogicalTypes}.
+   *
+   * @param logicalType The logical type
+   * @param utcTimezone whether to use UTC timezone for timestamp data type
+   */
+  public static Function<Object, Object> orderingValFunc(LogicalType 
logicalType, boolean utcTimezone) {
+    switch (logicalType.getTypeRoot()) {
+      case NULL:
+        return fieldVal -> null;
+      case TINYINT:
+        return fieldVal -> ((Byte) fieldVal).intValue();
+      case SMALLINT:
+        return fieldVal -> ((Short) fieldVal).intValue();
+      case DATE:
+        return fieldVal -> LocalDate.ofEpochDay((Long) fieldVal);
+      case CHAR:
+      case VARCHAR:
+        return Object::toString;
+      case BINARY:
+      case VARBINARY:
+        return fieldVal -> ByteBuffer.wrap((byte[]) fieldVal);
+      case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+        int precision1 = precision(logicalType);
+        if (precision1 <= 3) {
+          return fieldVal -> ((TimestampData) 
fieldVal).toInstant().toEpochMilli();
+        } else if (precision1 <= 6) {
+          return fieldVal -> {
+            Instant instant = ((TimestampData) fieldVal).toInstant();
+            return Math.addExact(Math.multiplyExact(instant.getEpochSecond(), 
1000_000), instant.getNano() / 1000);
+          };
+        } else {
+          throw new UnsupportedOperationException("Unsupported timestamp 
precision: " + precision1);
+        }
+      case TIMESTAMP_WITHOUT_TIME_ZONE:
+        int precision2 = precision(logicalType);
+        if (precision2 <= 3) {
+          return fieldVal -> utcTimezone ? ((TimestampData) 
fieldVal).toInstant().toEpochMilli() : ((TimestampData) 
fieldVal).toTimestamp().getTime();
+        } else if (precision2 <= 6) {
+          return fieldVal -> {
+            Instant instant = utcTimezone ? ((TimestampData) 
fieldVal).toInstant() : ((TimestampData) fieldVal).toTimestamp().toInstant();
+            return  Math.addExact(Math.multiplyExact(instant.getEpochSecond(), 
1000_000), instant.getNano() / 1000);
+          };
+        } else {
+          throw new UnsupportedOperationException("Unsupported timestamp 
precision: " + precision2);
+        }
+      case DECIMAL:
+        return fieldVal -> ((DecimalData) fieldVal).toBigDecimal();
+      default:
+        return fieldVal -> {
+          if (fieldVal == null) {
+            throw new HoodieValidationException("Ordering value(legacy as 
preCombine field value) can not be null");
+          }
+          return fieldVal;
+        };
+    }
+  }
+
+  /**
+   * Returns the precision of the given TIMESTAMP type.
+   */
+  public static int precision(LogicalType logicalType) {
+    if (logicalType instanceof TimestampType) {
+      return ((TimestampType) logicalType).getPrecision();
+    } else if (logicalType instanceof LocalZonedTimestampType) {
+      return ((LocalZonedTimestampType) logicalType).getPrecision();
+    } else {
+      throw new AssertionError("Unexpected type: " + logicalType);
+    }
+  }
+}
\ No newline at end of file
diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/common/model/HoodieSparkRecord.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/common/model/HoodieSparkRecord.java
index a7d78f5496f..76cbeec3e8c 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/common/model/HoodieSparkRecord.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/common/model/HoodieSparkRecord.java
@@ -185,6 +185,11 @@ public class HoodieSparkRecord extends 
HoodieRecord<InternalRow> {
     return objects;
   }
 
+  @Override
+  public Object getColumnValueAsJava(Schema recordSchema, String column, 
Properties props) {
+    throw new UnsupportedOperationException("Unsupported yet for " + 
this.getClass().getSimpleName());
+  }
+
   @Override
   public HoodieRecord joinWith(HoodieRecord other, Schema targetSchema) {
     StructType targetStructType = 
HoodieInternalRowUtils.getCachedSchema(targetSchema);
@@ -299,7 +304,12 @@ public class HoodieSparkRecord extends 
HoodieRecord<InternalRow> {
   }
 
   @Override
-  public Option<HoodieAvroIndexedRecord> toIndexedRecord(Schema recordSchema, 
Properties prop) throws IOException {
+  public Option<HoodieAvroIndexedRecord> toIndexedRecord(Schema recordSchema, 
Properties prop) {
+    throw new UnsupportedOperationException();
+  }
+
+  @Override
+  public byte[] getAvroBytes(Schema recordSchema, Properties props) throws 
IOException {
     throw new UnsupportedOperationException();
   }
 
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java 
b/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java
index 188311a5666..41a3f779d7b 100644
--- a/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java
+++ b/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java
@@ -174,6 +174,13 @@ public class HoodieAvroUtils {
     return 
Option.of(HoodieAvroUtils.indexedRecordToBytes(record.toIndexedRecord(schema, 
new Properties()).get().getData()));
   }
 
+  /**
+   * Convert a given avro record to bytes.
+   */
+  public static byte[] avroToBytes(IndexedRecord record) {
+    return indexedRecordToBytes(record);
+  }
+
   /**
    * Convert a given avro record to bytes.
    */
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
index 35e0d29be1d..ec2e9301c75 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
@@ -160,7 +160,7 @@ public class HoodieStorageConfig extends HoodieConfig {
       .withDocumentation("Control whether to write bloom filter or not. 
Default true. "
           + "We can set to false in non bloom index cases for CPU resource 
saving.");
 
-  public static final ConfigProperty<Boolean> PARQUET_WRITE_UTC_TIMEZONE = 
ConfigProperty
+  public static final ConfigProperty<Boolean> WRITE_UTC_TIMEZONE = 
ConfigProperty
       .key("hoodie.parquet.write.utc-timezone.enabled")
       .defaultValue(true)
       .markAdvanced()
@@ -515,7 +515,7 @@ public class HoodieStorageConfig extends HoodieConfig {
     }
 
     public Builder withWriteUtcTimezone(boolean writeUtcTimezone) {
-      storageConfig.setValue(PARQUET_WRITE_UTC_TIMEZONE, 
String.valueOf(writeUtcTimezone));
+      storageConfig.setValue(WRITE_UTC_TIMEZONE, 
String.valueOf(writeUtcTimezone));
       return this;
     }
 
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/model/BaseAvroPayload.java 
b/hudi-common/src/main/java/org/apache/hudi/common/model/BaseAvroPayload.java
index 85e46690287..4bea348c3b1 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/model/BaseAvroPayload.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/model/BaseAvroPayload.java
@@ -94,4 +94,8 @@ public abstract class BaseAvroPayload implements Serializable 
{
     Object deleteMarker = genericRecord.get(isDeleteKey);
     return (deleteMarker instanceof Boolean && (boolean) deleteMarker);
   }
+
+  public byte[] getRecordBytes() {
+    return recordBytes;
+  }
 }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java
index 9ae1e362841..368087c01be 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroIndexedRecord.java
@@ -114,6 +114,11 @@ public class HoodieAvroIndexedRecord extends 
HoodieRecord<IndexedRecord> {
     throw new UnsupportedOperationException();
   }
 
+  @Override
+  public Object getColumnValueAsJava(Schema recordSchema, String column, 
Properties props) {
+    throw new UnsupportedOperationException("Unsupported yet for " + 
this.getClass().getSimpleName());
+  }
+
   @Override
   public HoodieRecord joinWith(HoodieRecord other, Schema targetSchema) {
     GenericRecord record = HoodieAvroUtils.stitchRecords((GenericRecord) data, 
(GenericRecord) other.getData(), targetSchema);
@@ -212,6 +217,11 @@ public class HoodieAvroIndexedRecord extends 
HoodieRecord<IndexedRecord> {
     return Option.of(this);
   }
 
+  @Override
+  public byte[] getAvroBytes(Schema recordSchema, Properties props) {
+    return HoodieAvroUtils.avroToBytes(data);
+  }
+
   /**
    * NOTE: This method is declared final to make sure there's no polymorphism 
and therefore
    *       JIT compiler could perform more aggressive optimizations
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroRecord.java 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroRecord.java
index 285b292f6dd..3e2c5ce8d15 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroRecord.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieAvroRecord.java
@@ -120,6 +120,11 @@ public class HoodieAvroRecord<T extends 
HoodieRecordPayload> extends HoodieRecor
     return HoodieAvroUtils.getRecordColumnValues(this, columns, recordSchema, 
consistentLogicalTimestampEnabled);
   }
 
+  @Override
+  public Object getColumnValueAsJava(Schema recordSchema, String column, 
Properties props) {
+    throw new UnsupportedOperationException("Unsupported yet for " + 
this.getClass().getSimpleName());
+  }
+
   @Override
   public HoodieRecord joinWith(HoodieRecord other, Schema targetSchema) {
     throw new UnsupportedOperationException();
@@ -221,6 +226,16 @@ public class HoodieAvroRecord<T extends 
HoodieRecordPayload> extends HoodieRecor
     }
   }
 
+  @Override
+  public byte[] getAvroBytes(Schema recordSchema, Properties props) throws 
IOException {
+    if (data instanceof BaseAvroPayload) {
+      return ((BaseAvroPayload) getData()).getRecordBytes();
+    } else {
+      Option<IndexedRecord> avroData = getData().getInsertValue(recordSchema, 
props);
+      return avroData.map(HoodieAvroUtils::avroToBytes).orElse(new byte[0]);
+    }
+  }
+
   @Override
   protected final void writeRecordPayload(T payload, Kryo kryo, Output output) 
{
     // NOTE: Since [[orderingVal]] is polymorphic we have to write out its 
class
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieEmptyRecord.java 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieEmptyRecord.java
index 8183bde74d2..a4e1e44b758 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieEmptyRecord.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieEmptyRecord.java
@@ -97,6 +97,11 @@ public class HoodieEmptyRecord<T> extends HoodieRecord<T> {
     throw new UnsupportedOperationException();
   }
 
+  @Override
+  public Object getColumnValueAsJava(Schema recordSchema, String column, 
Properties props) {
+    throw new UnsupportedOperationException("Unsupported yet for " + 
this.getClass().getSimpleName());
+  }
+
   @Override
   public HoodieRecord joinWith(HoodieRecord other, Schema targetSchema) {
     throw new UnsupportedOperationException();
@@ -149,6 +154,11 @@ public class HoodieEmptyRecord<T> extends HoodieRecord<T> {
     return Option.empty();
   }
 
+  @Override
+  public byte[] getAvroBytes(Schema recordSchema, Properties props) {
+    return new byte[0];
+  }
+
   @Override
   public Option<Map<String, String>> getMetadata() {
     return Option.empty();
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecord.java 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecord.java
index f53a5874cfa..01062496d7e 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecord.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecord.java
@@ -393,6 +393,12 @@ public abstract class HoodieRecord<T> implements 
HoodieRecordCompatibilityInterf
    */
   public abstract Object[] getColumnValues(Schema recordSchema, String[] 
columns, boolean consistentLogicalTimestampEnabled);
 
+  /**
+   * Get column in record to support RDDCustomColumnsSortPartitioner
+   * @return column value
+   */
+  public abstract Object getColumnValueAsJava(Schema recordSchema, String 
column, Properties props);
+
   /**
    * Support bootstrap.
    */
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordCompatibilityInterface.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordCompatibilityInterface.java
index dd76fe9d3b5..1687fe20406 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordCompatibilityInterface.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordCompatibilityInterface.java
@@ -52,4 +52,6 @@ public interface HoodieRecordCompatibilityInterface {
   HoodieRecord truncateRecordKey(Schema recordSchema, Properties props, String 
keyFieldName) throws IOException;
 
   Option<HoodieAvroIndexedRecord> toIndexedRecord(Schema recordSchema, 
Properties props) throws IOException;
+
+  byte[] getAvroBytes(Schema recordSchema, Properties props) throws 
IOException;
 }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieAvroDataBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieAvroDataBlock.java
index f7a47bd4098..48e13894101 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieAvroDataBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieAvroDataBlock.java
@@ -25,6 +25,7 @@ import org.apache.hudi.common.fs.SizeAwareDataInputStream;
 import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType;
+import org.apache.hudi.common.util.CollectionUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.ClosableIterator;
 import org.apache.hudi.common.util.collection.CloseableMappingIterator;
@@ -33,13 +34,13 @@ import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.internal.schema.InternalSchema;
 import org.apache.hudi.io.SeekableDataInputStream;
 import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StorageConfiguration;
 
 import org.apache.avro.Schema;
 import org.apache.avro.generic.GenericDatumReader;
 import org.apache.avro.generic.GenericDatumWriter;
 import org.apache.avro.generic.IndexedRecord;
 import org.apache.avro.io.BinaryDecoder;
-import org.apache.avro.io.BinaryEncoder;
 import org.apache.avro.io.Decoder;
 import org.apache.avro.io.DecoderFactory;
 import org.apache.avro.io.Encoder;
@@ -77,8 +78,6 @@ import static 
org.apache.hudi.common.util.ValidationUtils.checkState;
  */
 public class HoodieAvroDataBlock extends HoodieDataBlock {
 
-  private final ThreadLocal<BinaryEncoder> encoderCache = new ThreadLocal<>();
-
   public HoodieAvroDataBlock(Supplier<SeekableDataInputStream> 
inputStreamSupplier,
                              Option<byte[]> content,
                              boolean readBlockLazily,
@@ -102,9 +101,8 @@ public class HoodieAvroDataBlock extends HoodieDataBlock {
   }
 
   @Override
-  protected byte[] serializeRecords(List<HoodieRecord> records, HoodieStorage 
storage) throws IOException {
+  protected ByteArrayOutputStream serializeRecords(List<HoodieRecord> records, 
HoodieStorage storage) throws IOException {
     Schema schema = AvroSchemaCache.intern(new 
Schema.Parser().parse(super.getLogBlockHeader().get(HeaderMetadataType.SCHEMA)));
-    GenericDatumWriter<IndexedRecord> writer = new 
GenericDatumWriter<>(schema);
     ByteArrayOutputStream baos = new ByteArrayOutputStream();
     try (DataOutputStream output = new DataOutputStream(baos)) {
       // 1. Write out the log block version
@@ -114,28 +112,22 @@ public class HoodieAvroDataBlock extends HoodieDataBlock {
       output.writeInt(records.size());
 
       // 3. Write the records
+      Properties props = initProperties(storage.getConf());
       for (HoodieRecord<?> s : records) {
-        ByteArrayOutputStream temp = new ByteArrayOutputStream();
-        BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(temp, 
encoderCache.get());
-        encoderCache.set(encoder);
         try {
           // Encode the record into bytes
           // Spark Record not support write avro log
-          IndexedRecord data = s.toIndexedRecord(schema, new 
Properties()).get().getData();
-          writer.write(data, encoder);
-          encoder.flush();
-
+          byte[] data = s.getAvroBytes(schema, props);
           // Write the record size
-          output.writeInt(temp.size());
+          output.writeInt(data.length);
           // Write the content
-          temp.writeTo(output);
+          output.write(data);
         } catch (IOException e) {
           throw new HoodieIOException("IOException converting 
HoodieAvroDataBlock to bytes", e);
         }
       }
-      encoderCache.remove();
     }
-    return baos.toByteArray();
+    return baos;
   }
 
   // TODO (na) - Break down content into smaller chunks of byte [] to be GC as 
they are used
@@ -396,6 +388,10 @@ public class HoodieAvroDataBlock extends HoodieDataBlock {
     }
   }
 
+  protected Properties initProperties(StorageConfiguration<?> storageConfig) {
+    return CollectionUtils.emptyProps();
+  }
+
   
//----------------------------------------------------------------------------------------
   //                                  DEPRECATED METHODS
   //
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCommandBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCommandBlock.java
index 2a900361d9c..965c57309b6 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCommandBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCommandBlock.java
@@ -22,6 +22,7 @@ import org.apache.hudi.common.util.Option;
 import org.apache.hudi.io.SeekableDataInputStream;
 import org.apache.hudi.storage.HoodieStorage;
 
+import java.io.ByteArrayOutputStream;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.function.Supplier;
@@ -62,7 +63,7 @@ public class HoodieCommandBlock extends HoodieLogBlock {
   }
 
   @Override
-  public byte[] getContentBytes(HoodieStorage storage) {
-    return new byte[0];
+  public ByteArrayOutputStream getContentBytes(HoodieStorage storage) {
+    return new ByteArrayOutputStream(0);
   }
 }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCorruptBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCorruptBlock.java
index 076ca31f271..2d2ee2967ca 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCorruptBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieCorruptBlock.java
@@ -22,6 +22,7 @@ import org.apache.hudi.common.util.Option;
 import org.apache.hudi.io.SeekableDataInputStream;
 import org.apache.hudi.storage.HoodieStorage;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.Map;
 import java.util.function.Supplier;
@@ -39,12 +40,12 @@ public class HoodieCorruptBlock extends HoodieLogBlock {
   }
 
   @Override
-  public byte[] getContentBytes(HoodieStorage storage) throws IOException {
+  public ByteArrayOutputStream getContentBytes(HoodieStorage storage) throws 
IOException {
     if (!getContent().isPresent() && readBlockLazily) {
       // read content from disk
       inflate();
     }
-    return getContent().get();
+    return getContentAsByteStream().get();
   }
 
   @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDataBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDataBlock.java
index 29dc17b052d..8abecae586c 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDataBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDataBlock.java
@@ -32,6 +32,7 @@ import org.apache.avro.Schema;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.HashSet;
 import java.util.Iterator;
@@ -116,14 +117,15 @@ public abstract class HoodieDataBlock extends 
HoodieLogBlock {
   }
 
   @Override
-  public byte[] getContentBytes(HoodieStorage storage) throws IOException {
+  public ByteArrayOutputStream getContentBytes(HoodieStorage storage) throws 
IOException {
     // In case this method is called before realizing records from content
     Option<byte[]> content = getContent();
 
     checkState(content.isPresent() || records.isPresent(), "Block is in 
invalid state");
 
-    if (content.isPresent()) {
-      return content.get();
+    Option<ByteArrayOutputStream> baosOpt = getContentAsByteStream();
+    if (baosOpt.isPresent()) {
+      return baosOpt.get();
     }
 
     return serializeRecords(records.get(), storage);
@@ -326,7 +328,7 @@ public abstract class HoodieDataBlock extends 
HoodieLogBlock {
     );
   }
 
-  protected abstract byte[] serializeRecords(List<HoodieRecord> records, 
HoodieStorage storage) throws IOException;
+  protected abstract ByteArrayOutputStream serializeRecords(List<HoodieRecord> 
records, HoodieStorage storage) throws IOException;
 
   protected abstract <T> ClosableIterator<HoodieRecord<T>> 
deserializeRecords(byte[] content, HoodieRecordType type) throws IOException;
 
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDeleteBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDeleteBlock.java
index fe863546af3..bb2db278d64 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDeleteBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieDeleteBlock.java
@@ -88,12 +88,11 @@ public class HoodieDeleteBlock extends HoodieLogBlock {
   }
 
   @Override
-  public byte[] getContentBytes(HoodieStorage storage) throws IOException {
-    Option<byte[]> content = getContent();
-
+  public ByteArrayOutputStream getContentBytes(HoodieStorage storage) throws 
IOException {
     // In case this method is called before realizing keys from content
-    if (content.isPresent()) {
-      return content.get();
+    Option<ByteArrayOutputStream> baosOpt = getContentAsByteStream();
+    if (baosOpt.isPresent()) {
+      return baosOpt.get();
     } else if (readBlockLazily && recordsToDelete == null) {
       // read block lazily
       getRecordsToDelete();
@@ -105,7 +104,7 @@ public class HoodieDeleteBlock extends HoodieLogBlock {
     byte[] bytesToWrite = (version <= 2) ? serializeV2() : serializeV3();
     output.writeInt(bytesToWrite.length);
     output.write(bytesToWrite);
-    return baos.toByteArray();
+    return baos;
   }
 
   public DeleteRecord[] getRecordsToDelete() {
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieHFileDataBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieHFileDataBlock.java
index a71bf203246..6480ce83eff 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieHFileDataBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieHFileDataBlock.java
@@ -41,6 +41,7 @@ import org.apache.avro.generic.IndexedRecord;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.Collections;
 import java.util.HashMap;
@@ -99,7 +100,7 @@ public class HoodieHFileDataBlock extends HoodieDataBlock {
   }
 
   @Override
-  protected byte[] serializeRecords(List<HoodieRecord> records, HoodieStorage 
storage) throws IOException {
+  protected ByteArrayOutputStream serializeRecords(List<HoodieRecord> records, 
HoodieStorage storage) throws IOException {
     Schema writerSchema = new Schema.Parser().parse(
         
super.getLogBlockHeader().get(HoodieLogBlock.HeaderMetadataType.SCHEMA));
     return 
HoodieIOFactory.getIOFactory(storage).getFileFormatUtils(HoodieFileFormat.HFILE)
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieLogBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieLogBlock.java
index 78c39ff5faf..ad05c10854c 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieLogBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieLogBlock.java
@@ -93,7 +93,7 @@ public abstract class HoodieLogBlock {
   }
 
   // Return the bytes representation of the data belonging to a LogBlock
-  public byte[] getContentBytes(HoodieStorage storage) throws IOException {
+  public ByteArrayOutputStream getContentBytes(HoodieStorage storage) throws 
IOException {
     throw new HoodieException("No implementation was provided");
   }
 
@@ -345,6 +345,21 @@ public abstract class HoodieLogBlock {
     return Option.of(content);
   }
 
+  /**
+   * Return bytes content as a {@link ByteArrayOutputStream}.
+   *
+   * @return a {@link ByteArrayOutputStream} contains the block content bytes
+   */
+  protected Option<ByteArrayOutputStream> getContentAsByteStream() throws 
IOException {
+    if (content.isEmpty()) {
+      return Option.empty();
+    }
+    byte[] contentBytes = content.get();
+    ByteArrayOutputStream baos = new 
ByteArrayOutputStream(contentBytes.length);
+    baos.write(contentBytes);
+    return Option.of(baos);
+  }
+
   protected Supplier<SeekableDataInputStream> getInputStreamSupplier() {
     return inputStreamSupplier;
   }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieParquetDataBlock.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieParquetDataBlock.java
index 0c47a8092a9..424703a9f3c 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieParquetDataBlock.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/block/HoodieParquetDataBlock.java
@@ -33,6 +33,7 @@ import org.apache.hudi.storage.inline.InLineFSUtils;
 
 import org.apache.avro.Schema;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.HashMap;
 import java.util.List;
@@ -88,7 +89,7 @@ public class HoodieParquetDataBlock extends HoodieDataBlock {
   }
 
   @Override
-  protected byte[] serializeRecords(List<HoodieRecord> records, HoodieStorage 
storage) throws IOException {
+  protected ByteArrayOutputStream serializeRecords(List<HoodieRecord> records, 
HoodieStorage storage) throws IOException {
     Map<String, String> paramsMap = new HashMap<>();
     paramsMap.put(PARQUET_COMPRESSION_CODEC_NAME.key(), 
compressionCodecName.get());
     paramsMap.put(PARQUET_COMPRESSION_RATIO_FRACTION.key(), 
String.valueOf(expectedCompressionRatio.get()));
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/util/FileFormatUtils.java 
b/hudi-common/src/main/java/org/apache/hudi/common/util/FileFormatUtils.java
index c65313703eb..3216561e861 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/util/FileFormatUtils.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/util/FileFormatUtils.java
@@ -39,6 +39,7 @@ import org.apache.avro.generic.GenericRecord;
 
 import javax.annotation.Nonnull;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.HashMap;
 import java.util.HashSet;
@@ -319,11 +320,11 @@ public abstract class FileFormatUtils {
    * @return byte array after serialization.
    * @throws IOException upon serialization error.
    */
-  public abstract byte[] serializeRecordsToLogBlock(HoodieStorage storage,
-                                                    List<HoodieRecord> records,
-                                                    Schema writerSchema,
-                                                    Schema readerSchema, 
String keyFieldName,
-                                                    Map<String, String> 
paramsMap) throws IOException;
+  public abstract ByteArrayOutputStream 
serializeRecordsToLogBlock(HoodieStorage storage,
+                                                                   
List<HoodieRecord> records,
+                                                                   Schema 
writerSchema,
+                                                                   Schema 
readerSchema, String keyFieldName,
+                                                                   Map<String, 
String> paramsMap) throws IOException;
 
   /**
    * Serializes Hudi records to the log block and collect column range 
metadata.
@@ -337,13 +338,13 @@ public abstract class FileFormatUtils {
    * @return pair of byte array after serialization and format metadata.
    * @throws IOException upon serialization error.
    */
-  public abstract Pair<byte[], Object> 
serializeRecordsToLogBlock(HoodieStorage storage,
-                                                                  
Iterator<HoodieRecord> records,
-                                                                  
HoodieRecord.HoodieRecordType recordType,
-                                                                  Schema 
writerSchema,
-                                                                  Schema 
readerSchema,
-                                                                  String 
keyFieldName,
-                                                                  Map<String, 
String> paramsMap) throws IOException;
+  public abstract Pair<ByteArrayOutputStream, Object> 
serializeRecordsToLogBlock(HoodieStorage storage,
+                                                                               
  Iterator<HoodieRecord> records,
+                                                                               
  HoodieRecord.HoodieRecordType recordType,
+                                                                               
  Schema writerSchema,
+                                                                               
  Schema readerSchema,
+                                                                               
  String keyFieldName,
+                                                                               
  Map<String, String> paramsMap) throws IOException;
 
   // -------------------------------------------------------------------------
   //  Inner Class
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
index 27e7c727020..52e7bf948b9 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
@@ -39,6 +39,7 @@ import org.apache.hudi.avro.model.TimestampMicrosWrapper;
 import org.apache.hudi.common.bloom.BloomFilter;
 import org.apache.hudi.common.config.HoodieConfig;
 import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.data.HoodieAccumulator;
 import org.apache.hudi.common.data.HoodieAtomicLongAccumulator;
 import org.apache.hudi.common.data.HoodieData;
@@ -136,6 +137,7 @@ import java.util.LinkedList;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Properties;
 import java.util.Set;
 import java.util.TreeMap;
 import java.util.UUID;
@@ -241,7 +243,11 @@ public class HoodieTableMetadataUtil {
    *         the collection of provided records
    */
   public static Map<String, HoodieColumnRangeMetadata<Comparable>> 
collectColumnRangeMetadata(
-      List<HoodieRecord> records, List<Pair<String, Schema.Field>> 
targetFields, String filePath, Schema recordSchema) {
+      List<HoodieRecord> records,
+      List<Pair<String, Schema.Field>> targetFields,
+      String filePath,
+      Schema recordSchema,
+      StorageConfiguration<?> storageConfig) {
     // Helper class to calculate column stats
     class ColumnStats {
       Object minValue;
@@ -252,6 +258,9 @@ public class HoodieTableMetadataUtil {
 
     HashMap<String, ColumnStats> allColumnStats = new HashMap<>();
 
+    final Properties properties = new Properties();
+    properties.setProperty(HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(),
+        storageConfig.getString(HoodieStorageConfig.WRITE_UTC_TIMEZONE.key(), 
HoodieStorageConfig.WRITE_UTC_TIMEZONE.defaultValue().toString()));
     // Collect stats for all columns by iterating through records while 
accounting
     // corresponding stats
     records.forEach((record) -> {
@@ -273,6 +282,8 @@ public class HoodieTableMetadataUtil {
           if (fieldSchema.getType() == Schema.Type.INT && 
fieldSchema.getLogicalType() != null && fieldSchema.getLogicalType() == 
LogicalTypes.date()) {
             fieldValue = java.sql.Date.valueOf(LocalDate.ofEpochDay((Integer) 
fieldValue).toString());
           }
+        } else if (record.getRecordType() == HoodieRecordType.FLINK) {
+          fieldValue = record.getColumnValueAsJava(recordSchema, fieldName, 
properties);
         } else {
           throw new HoodieException(String.format("Unknown record type: %s", 
record.getRecordType()));
         }
@@ -1713,7 +1724,7 @@ public class HoodieTableMetadataUtil {
         return Collections.emptyList();
       }
       Map<String, HoodieColumnRangeMetadata<Comparable>> 
columnRangeMetadataMap =
-          collectColumnRangeMetadata(records, fieldsToIndex, 
getFileNameFromPath(filePath), writerSchemaOpt.get());
+          collectColumnRangeMetadata(records, fieldsToIndex, 
getFileNameFromPath(filePath), writerSchemaOpt.get(), 
datasetMetaClient.getStorage().getConf());
       return new ArrayList<>(columnRangeMetadataMap.values());
     }
     return Collections.emptyList();
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieAvroDataBlock.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieAvroDataBlock.java
index dae7e6e841c..716f08a88b3 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieAvroDataBlock.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieAvroDataBlock.java
@@ -32,6 +32,9 @@ import 
org.apache.hudi.common.util.io.ByteBufferBackedInputStream;
 import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.io.ByteArraySeekableDataInputStream;
 import org.apache.hudi.io.SeekableDataInputStream;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StorageConfiguration;
+
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.CsvSource;
@@ -50,6 +53,8 @@ import java.util.stream.Collectors;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
 
 public class TestHoodieAvroDataBlock {
   // Record key field of the test data
@@ -304,7 +309,9 @@ public class TestHoodieAvroDataBlock {
     Map<HeaderMetadataType, String> header = new HashMap<>();
     header.put(HeaderMetadataType.SCHEMA, schema.toString());
 
-    return new HoodieAvroDataBlock(records, header, 
RECORD_KEY_FIELD).getContentBytes(null);
+    HoodieStorage storage = mock(HoodieStorage.class);
+    when(storage.getConf()).thenReturn(mock(StorageConfiguration.class));
+    return new HoodieAvroDataBlock(records, header, 
RECORD_KEY_FIELD).getContentBytes(storage).toByteArray();
   }
 
   /**
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/RowDataStreamWriteFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/RowDataStreamWriteFunction.java
index 4f69dfe91c0..84f38a80145 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/RowDataStreamWriteFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/RowDataStreamWriteFunction.java
@@ -219,7 +219,7 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
   }
 
   private void initRecordConverter() {
-    this.recordConverter = RecordConverter.getInstance(config, rowType, 
keyGen, writeClient.getConfig(), metaClient.getTableConfig());
+    this.recordConverter = RecordConverter.getInstance(keyGen);
   }
 
   private void initMergeClass() {
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/transform/RecordConverter.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/transform/RecordConverter.java
index 86095730376..ec0c7f68fc2 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/transform/RecordConverter.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/transform/RecordConverter.java
@@ -18,28 +18,14 @@
 
 package org.apache.hudi.sink.transform;
 
-import org.apache.hudi.client.model.HoodieFlinkAvroRecord;
 import org.apache.hudi.client.model.HoodieFlinkRecord;
 import org.apache.hudi.common.model.HoodieKey;
 import org.apache.hudi.common.model.HoodieOperation;
 import org.apache.hudi.common.model.HoodieRecord;
-import org.apache.hudi.common.table.HoodieTableConfig;
-import org.apache.hudi.common.table.log.block.HoodieLogBlock;
-import org.apache.hudi.config.HoodieWriteConfig;
-import org.apache.hudi.configuration.FlinkOptions;
-import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.sink.bulk.RowDataKeyGen;
 import org.apache.hudi.table.action.commit.BucketInfo;
-import org.apache.hudi.util.CommonClientUtils;
-import org.apache.hudi.util.OrderingValueExtractor;
-import org.apache.hudi.util.RowDataToAvroConverters;
-import org.apache.hudi.util.StreamerUtil;
 
-import org.apache.avro.Schema;
-import org.apache.avro.generic.GenericRecord;
-import org.apache.flink.configuration.Configuration;
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.types.logical.RowType;
 
 import java.io.Serializable;
 
@@ -49,41 +35,12 @@ import java.io.Serializable;
 public interface RecordConverter extends Serializable {
   HoodieRecord convert(RowData dataRow, BucketInfo bucketInfo);
 
-  static RecordConverter getInstance(
-      Configuration flinkConf,
-      RowType rowType,
-      RowDataKeyGen keyGen,
-      HoodieWriteConfig writeConfig,
-      HoodieTableConfig tableConfig) {
-    // construct flink record according to the log block format type
-    HoodieLogBlock.HoodieLogBlockType logBlockType = 
CommonClientUtils.getLogBlockType(writeConfig, tableConfig);
-    OrderingValueExtractor orderingValueExtractor = 
OrderingValueExtractor.getInstance(flinkConf, rowType);
-    if (logBlockType == HoodieLogBlock.HoodieLogBlockType.PARQUET_DATA_BLOCK) {
-      return (dataRow, bucketInfo) -> {
-        String key = keyGen.getRecordKey(dataRow);
-        Comparable<?> orderingValue = 
orderingValueExtractor.getOrderingValue(dataRow);
-        HoodieOperation operation = 
HoodieOperation.fromValue(dataRow.getRowKind().toByteValue());
-        HoodieKey hoodieKey = new HoodieKey(key, 
bucketInfo.getPartitionPath());
-        return new HoodieFlinkRecord(hoodieKey, operation, orderingValue, 
dataRow);
-      };
-    } else if (logBlockType == 
HoodieLogBlock.HoodieLogBlockType.AVRO_DATA_BLOCK) {
-      return new RecordConverter() {
-        private final Schema avroSchema = 
StreamerUtil.getSourceSchema(flinkConf);
-        private final RowDataToAvroConverters.RowDataToAvroConverter converter 
= RowDataToAvroConverters.createConverter(rowType, 
flinkConf.get(FlinkOptions.WRITE_UTC_TIMEZONE));
-
-        @Override
-        public HoodieRecord convert(RowData dataRow, BucketInfo bucketInfo) {
-          String key = keyGen.getRecordKey(dataRow);
-          Comparable<?> orderingValue = 
orderingValueExtractor.getOrderingValue(dataRow);
-          HoodieOperation operation = 
HoodieOperation.fromValue(dataRow.getRowKind().toByteValue());
-          HoodieKey hoodieKey = new HoodieKey(key, 
bucketInfo.getPartitionPath());
-
-          GenericRecord record = (GenericRecord) converter.convert(avroSchema, 
dataRow);
-          return new HoodieFlinkAvroRecord(hoodieKey, operation, 
orderingValue, record);
-        }
-      };
-    } else {
-      throw new HoodieException("Unsupported log block type: " + logBlockType);
-    }
+  static RecordConverter getInstance(RowDataKeyGen keyGen) {
+    return (dataRow, bucketInfo) -> {
+      String key = keyGen.getRecordKey(dataRow);
+      HoodieOperation operation = 
HoodieOperation.fromValue(dataRow.getRowKind().toByteValue());
+      HoodieKey hoodieKey = new HoodieKey(key, bucketInfo.getPartitionPath());
+      return new HoodieFlinkRecord(hoodieKey, operation, dataRow);
+    };
   }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/OrderingValueExtractor.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/OrderingValueExtractor.java
deleted file mode 100644
index fdebb738a9b..00000000000
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/OrderingValueExtractor.java
+++ /dev/null
@@ -1,57 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *      http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.hudi.util;
-
-import org.apache.hudi.common.model.HoodieRecord;
-import org.apache.hudi.common.model.WriteOperationType;
-import org.apache.hudi.common.util.StringUtils;
-import org.apache.hudi.configuration.FlinkOptions;
-import org.apache.hudi.configuration.OptionsResolver;
-
-import org.apache.flink.configuration.Configuration;
-import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.types.logical.LogicalType;
-import org.apache.flink.table.types.logical.RowType;
-
-import java.io.Serializable;
-
-/**
- * Interface for extracting ordering value from RowData.
- */
-public interface OrderingValueExtractor extends Serializable {
-  Comparable<?> getOrderingValue(RowData rowData);
-
-  static OrderingValueExtractor getInstance(Configuration conf, RowType 
rowType) {
-    String fieldName = OptionsResolver.getPreCombineField(conf);
-    boolean needCombine = conf.getBoolean(FlinkOptions.PRE_COMBINE)
-        || 
WriteOperationType.fromValue(conf.getString(FlinkOptions.OPERATION)) == 
WriteOperationType.UPSERT;
-    boolean shouldCombine = needCombine && 
!StringUtils.isNullOrEmpty(fieldName);
-    if (!shouldCombine) {
-      // returns a natual order value extractor.
-      return rowData -> HoodieRecord.DEFAULT_ORDERING_VALUE;
-    }
-    final int fieldPos = rowType.getFieldNames().indexOf(fieldName);
-    final LogicalType fieldType = rowType.getTypeAt(fieldPos);
-    final RowData.FieldGetter orderingValueFieldGetter = 
RowData.createFieldGetter(fieldType, fieldPos);
-    boolean utcTimezone = conf.get(FlinkOptions.WRITE_UTC_TIMEZONE);
-
-    return rowData -> (Comparable<?>) DataTypeUtils.resolveOrderingValue(
-        fieldType, orderingValueFieldGetter.getFieldOrNull(rowData), 
utcTimezone);
-  }
-}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
index 7426ffc14c4..d1ca7032ea1 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
@@ -2646,19 +2646,6 @@ public class ITTestHoodieDataSource {
     return Stream.of(data).map(Arguments::of);
   }
 
-  /**
-   * Return test params => (HoodieTableType, LogBlockType).
-   */
-  private static Stream<Arguments> tableTypeAndLogBlockTypeParams() {
-    Object[][] data =
-        new Object[][] {
-            {HoodieTableType.COPY_ON_WRITE, "avro"},
-            {HoodieTableType.COPY_ON_WRITE, "parquet"},
-            {HoodieTableType.MERGE_ON_READ, "avro"},
-            {HoodieTableType.MERGE_ON_READ, "parquet"}};
-    return Stream.of(data).map(Arguments::of);
-  }
-
   public static List<Arguments> testBulkInsertWithPartitionBucketIndexParams() 
{
     return asList(
         Arguments.of("bulk_insert", COPY_ON_WRITE.name()),
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/table/log/HoodieLogFormatWriter.java
 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/table/log/HoodieLogFormatWriter.java
index e2c75aa9367..f949c26b268 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/table/log/HoodieLogFormatWriter.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/table/log/HoodieLogFormatWriter.java
@@ -34,6 +34,7 @@ import org.apache.hadoop.ipc.RemoteException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.Collections;
 import java.util.List;
@@ -151,12 +152,12 @@ public class HoodieLogFormatWriter implements 
HoodieLogFormat.Writer {
       // bytes for header
       byte[] headerBytes = 
HoodieLogBlock.getHeaderMetadataBytes(block.getLogBlockHeader());
       // content bytes
-      byte[] content = block.getContentBytes(storage);
+      ByteArrayOutputStream content = block.getContentBytes(storage);
       // bytes for footer
       byte[] footerBytes = 
HoodieLogBlock.getFooterMetadataBytes(block.getLogBlockFooter());
 
       // 2. Write the total size of the block (excluding Magic)
-      outputStream.writeLong(getLogBlockLength(content.length, 
headerBytes.length, footerBytes.length));
+      outputStream.writeLong(getLogBlockLength(content.size(), 
headerBytes.length, footerBytes.length));
 
       // 3. Write the version of this log block
       outputStream.writeInt(currentLogFormatVersion.getVersion());
@@ -166,9 +167,9 @@ public class HoodieLogFormatWriter implements 
HoodieLogFormat.Writer {
       // 5. Write the headers for the log block
       outputStream.write(headerBytes);
       // 6. Write the size of the content block
-      outputStream.writeLong(content.length);
+      outputStream.writeLong(content.size());
       // 7. Write the contents of the data block
-      outputStream.write(content);
+      content.writeTo(outputStream);
       // 8. Write the footers for the log block
       outputStream.write(footerBytes);
       // 9. Write the total size of the log block (including magic) which is 
everything written
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/HFileUtils.java 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/HFileUtils.java
index 314a7d0dcfe..001a4dd3474 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/HFileUtils.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/HFileUtils.java
@@ -161,7 +161,7 @@ public class HFileUtils extends FileFormatUtils {
   }
 
   @Override
-  public byte[] serializeRecordsToLogBlock(HoodieStorage storage,
+  public ByteArrayOutputStream serializeRecordsToLogBlock(HoodieStorage 
storage,
                                            List<HoodieRecord> records,
                                            Schema writerSchema,
                                            Schema readerSchema,
@@ -229,7 +229,7 @@ public class HFileUtils extends FileFormatUtils {
     ostream.flush();
     ostream.close();
 
-    return baos.toByteArray();
+    return baos;
   }
 
   /**
@@ -246,7 +246,7 @@ public class HFileUtils extends FileFormatUtils {
   }
 
   @Override
-  public Pair<byte[], Object> serializeRecordsToLogBlock(
+  public Pair<ByteArrayOutputStream, Object> serializeRecordsToLogBlock(
       HoodieStorage storage,
       Iterator<HoodieRecord> records,
       HoodieRecord.HoodieRecordType recordType,
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/OrcUtils.java 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/OrcUtils.java
index 8b7fae0e18d..5161e344a93 100644
--- a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/OrcUtils.java
+++ b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/OrcUtils.java
@@ -49,6 +49,7 @@ import org.apache.orc.RecordReader;
 import org.apache.orc.TypeDescription;
 import org.apache.orc.Writer;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
@@ -307,23 +308,23 @@ public class OrcUtils extends FileFormatUtils {
   }
 
   @Override
-  public byte[] serializeRecordsToLogBlock(HoodieStorage storage,
-                                           List<HoodieRecord> records,
-                                           Schema writerSchema,
-                                           Schema readerSchema,
-                                           String keyFieldName,
-                                           Map<String, String> paramsMap) 
throws IOException {
+  public ByteArrayOutputStream serializeRecordsToLogBlock(HoodieStorage 
storage,
+                                                          List<HoodieRecord> 
records,
+                                                          Schema writerSchema,
+                                                          Schema readerSchema,
+                                                          String keyFieldName,
+                                                          Map<String, String> 
paramsMap) throws IOException {
     throw new UnsupportedOperationException("Hudi log blocks do not support 
ORC format yet");
   }
 
   @Override
-  public Pair<byte[], Object> serializeRecordsToLogBlock(HoodieStorage storage,
-                                                         
Iterator<HoodieRecord> records,
-                                                         
HoodieRecord.HoodieRecordType recordType,
-                                                         Schema writerSchema,
-                                                         Schema readerSchema,
-                                                         String keyFieldName,
-                                                         Map<String, String> 
paramsMap) throws IOException {
+  public Pair<ByteArrayOutputStream, Object> 
serializeRecordsToLogBlock(HoodieStorage storage,
+                                                                        
Iterator<HoodieRecord> records,
+                                                                        
HoodieRecord.HoodieRecordType recordType,
+                                                                        Schema 
writerSchema,
+                                                                        Schema 
readerSchema,
+                                                                        String 
keyFieldName,
+                                                                        
Map<String, String> paramsMap) throws IOException {
     throw new UnsupportedOperationException("Hudi log blocks do not support 
ORC format yet");
   }
 }
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
index bddc15b2893..586b8367b32 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
@@ -382,14 +382,14 @@ public class ParquetUtils extends FileFormatUtils {
   }
 
   @Override
-  public byte[] serializeRecordsToLogBlock(HoodieStorage storage,
+  public ByteArrayOutputStream serializeRecordsToLogBlock(HoodieStorage 
storage,
                                            List<HoodieRecord> records,
                                            Schema writerSchema,
                                            Schema readerSchema,
                                            String keyFieldName,
                                            Map<String, String> paramsMap) 
throws IOException {
     if (records.size() == 0) {
-      return new byte[0];
+      return new ByteArrayOutputStream(0);
     }
 
     ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
@@ -407,17 +407,17 @@ public class ParquetUtils extends FileFormatUtils {
       }
       outputStream.flush();
     }
-    return outputStream.toByteArray();
+    return outputStream;
   }
 
   @Override
-  public Pair<byte[], Object> serializeRecordsToLogBlock(HoodieStorage storage,
-                                                         
Iterator<HoodieRecord> recordItr,
-                                                         
HoodieRecord.HoodieRecordType recordType,
-                                                         Schema writerSchema,
-                                                         Schema readerSchema,
-                                                         String keyFieldName,
-                                                         Map<String, String> 
paramsMap) throws IOException {
+  public Pair<ByteArrayOutputStream, Object> 
serializeRecordsToLogBlock(HoodieStorage storage,
+                                                                        
Iterator<HoodieRecord> recordItr,
+                                                                        
HoodieRecord.HoodieRecordType recordType,
+                                                                        Schema 
writerSchema,
+                                                                        Schema 
readerSchema,
+                                                                        String 
keyFieldName,
+                                                                        
Map<String, String> paramsMap) throws IOException {
     ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
     HoodieConfig config = new HoodieConfig();
     paramsMap.entrySet().stream().forEach(entry -> 
config.setValue(entry.getKey(), entry.getValue()));
@@ -434,7 +434,7 @@ public class ParquetUtils extends FileFormatUtils {
     }
     outputStream.flush();
     parquetWriter.close();
-    return Pair.of(outputStream.toByteArray(), 
parquetWriter.getFileFormatMetadata());
+    return Pair.of(outputStream, parquetWriter.getFileFormatMetadata());
   }
 
   static class RecordKeysFilterFunction implements Function<String, Boolean> {
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/functional/TestHoodieLogFormat.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/functional/TestHoodieLogFormat.java
index 25545de976c..352a4946f02 100755
--- 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/functional/TestHoodieLogFormat.java
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/functional/TestHoodieLogFormat.java
@@ -472,7 +472,7 @@ public class TestHoodieLogFormat extends 
HoodieCommonTestHarness {
     Map<HoodieLogBlock.HeaderMetadataType, String> header = new HashMap<>();
     header.put(HoodieLogBlock.HeaderMetadataType.INSTANT_TIME, "100");
     header.put(HoodieLogBlock.HeaderMetadataType.SCHEMA, 
getSimpleSchema().toString());
-    byte[] dataBlockContentBytes = getDataBlock(DEFAULT_DATA_BLOCK_TYPE, 
records, header).getContentBytes(storage);
+    byte[] dataBlockContentBytes = getDataBlock(DEFAULT_DATA_BLOCK_TYPE, 
records, header).getContentBytes(storage).toByteArray();
     HoodieLogBlock.HoodieLogBlockContentLocation logBlockContentLoc = new 
HoodieLogBlock.HoodieLogBlockContentLocation(
         HoodieTestUtils.getStorage(basePath), null, 0, 
dataBlockContentBytes.length, 0);
     HoodieDataBlock reusableDataBlock = new HoodieAvroDataBlock(null, 
Option.ofNullable(dataBlockContentBytes), false,
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieDeleteBlock.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieDeleteBlock.java
index ba513e391cd..35525ee9d09 100644
--- 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieDeleteBlock.java
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/table/log/block/TestHoodieDeleteBlock.java
@@ -122,7 +122,7 @@ public class TestHoodieDeleteBlock {
       deleteRecordList.add(Pair.of(dr, -1L));
     }
     HoodieDeleteBlock deleteBlock = new HoodieDeleteBlock(deleteRecordList, 
new HashMap<>());
-    byte[] contentBytes = 
deleteBlock.getContentBytes(HoodieTestUtils.getDefaultStorage());
+    byte[] contentBytes = 
deleteBlock.getContentBytes(HoodieTestUtils.getDefaultStorage()).toByteArray();
     HoodieDeleteBlock deserializeDeleteBlock = new HoodieDeleteBlock(
         Option.of(contentBytes), null, true, Option.empty(), new HashMap<>(), 
new HashMap<>());
     DeleteRecord[] deserializedDeleteRecords = 
deserializeDeleteBlock.getRecordsToDelete();
diff --git 
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecord.java 
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecord.java
index 2f5436e41f5..d17b8eca6c1 100644
--- a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecord.java
+++ b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecord.java
@@ -143,6 +143,11 @@ public class HoodieHiveRecord extends 
HoodieRecord<ArrayWritable> {
     return objects;
   }
 
+  @Override
+  public Object getColumnValueAsJava(Schema recordSchema, String column, 
Properties props) {
+    throw new UnsupportedOperationException("Unsupported yet for " + 
this.getClass().getSimpleName());
+  }
+
   @Override
   public HoodieRecord joinWith(HoodieRecord other, Schema targetSchema) {
     throw new UnsupportedOperationException("Not supported for 
HoodieHiveRecord");
@@ -212,6 +217,11 @@ public class HoodieHiveRecord extends 
HoodieRecord<ArrayWritable> {
     throw new UnsupportedOperationException("Not supported for 
HoodieHiveRecord");
   }
 
+  @Override
+  public byte[] getAvroBytes(Schema recordSchema, Properties props) throws 
IOException {
+    throw new UnsupportedOperationException("Not supported for 
HoodieHiveRecord");
+  }
+
   private Object getValue(String name) {
     return HoodieArrayWritableAvroUtils.getWritableValue(data, 
objectInspector, name);
   }
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/test/java/org/apache/hudi/testutils/LogFileColStatsTestUtil.java
 
b/hudi-spark-datasource/hudi-spark-common/src/test/java/org/apache/hudi/testutils/LogFileColStatsTestUtil.java
index 242f81dcefc..519a8811eb5 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/test/java/org/apache/hudi/testutils/LogFileColStatsTestUtil.java
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/test/java/org/apache/hudi/testutils/LogFileColStatsTestUtil.java
@@ -68,7 +68,7 @@ public class LogFileColStatsTestUtil {
         return Option.empty();
       }
       Map<String, HoodieColumnRangeMetadata<Comparable>> 
columnRangeMetadataMap =
-          collectColumnRangeMetadata(records, fieldsToIndex, filePath, 
writerSchemaOpt.get());
+          collectColumnRangeMetadata(records, fieldsToIndex, filePath, 
writerSchemaOpt.get(), datasetMetaClient.getStorageConf());
       List<HoodieColumnRangeMetadata<Comparable>> columnRangeMetadataList = 
new ArrayList<>(columnRangeMetadataMap.values());
       return Option.of(getColStatsEntry(filePath, columnRangeMetadataList));
     } else {

Reply via email to