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 66f216573c4 [HUDI-9190] RowDataLogWriteHandle supports writing Avro 
data block (#13142)
66f216573c4 is described below

commit 66f216573c4383a597772812479efa47f1d1ef32
Author: Shuo Cheng <[email protected]>
AuthorDate: Thu Apr 17 13:08:14 2025 +0800

    [HUDI-9190] RowDataLogWriteHandle supports writing Avro data block (#13142)
    
    * introduce HoodieFlinkAvroRecord
    * add record converter
    * add ordering value resolver
    
    ---------
    
    Co-authored-by: danny0405 <[email protected]>
---
 .../org/apache/hudi/io/HoodieAppendHandle.java     | 22 +-----
 .../java/org/apache/hudi/table/HoodieTable.java    |  7 +-
 .../org/apache/hudi/util/CommonClientUtils.java    | 35 +++++++++
 ...FlinkRecord.java => HoodieFlinkAvroRecord.java} | 78 ++++++++++---------
 .../hudi/client/model/HoodieFlinkRecord.java       |  8 +-
 .../java/org/apache/hudi/io/FlinkAppendHandle.java |  7 +-
 .../apache/hudi/io/FlinkWriteHandleFactory.java    |  4 +-
 .../hudi/io/v2/FlinkRowDataHandleFactory.java      | 33 +++++---
 .../apache/hudi/io/v2/RowDataLogWriteHandle.java   | 82 ++++++--------------
 .../hudi/table/HoodieFlinkMergeOnReadTable.java    |  8 +-
 .../RowDataUpsertDeltaCommitActionExecutor.java    |  7 +-
 .../apache/hudi/util/OrderingValueExtractor.java   | 30 --------
 .../table/log/block/HoodieAvroDataBlock.java       |  6 +-
 .../hudi/sink/RowDataStreamWriteFunction.java      | 90 +++++++---------------
 ...RowDataConsistentBucketStreamWriteFunction.java |  2 +-
 .../hudi/sink/transform/RecordConverter.java       | 89 +++++++++++++++++++++
 .../java/org/apache/hudi/util/DataTypeUtils.java   | 65 ++++++++++++++++
 .../apache/hudi/util/OrderingValueExtractor.java   | 57 ++++++++++++++
 .../apache/hudi/table/ITTestHoodieDataSource.java  | 77 ++++++++++++++++++
 19 files changed, 468 insertions(+), 239 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 a68f3629cf2..7571d173aaa 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
@@ -61,6 +61,7 @@ import org.apache.hudi.exception.HoodieUpsertException;
 import org.apache.hudi.metadata.HoodieTableMetadataUtil;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.util.CommonClientUtils;
 import org.apache.hudi.util.Lazy;
 
 import org.apache.avro.Schema;
@@ -499,7 +500,7 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
             ? HoodieRecord.RECORD_KEY_METADATA_FIELD
             : 
hoodieTable.getMetaClient().getTableConfig().getRecordKeyFieldProp();
 
-        dataBlock = getDataBlock(config, pickLogDataBlockFormat(), recordList,
+        dataBlock = getDataBlock(config, getLogBlockType(), recordList,
             getUpdatedHeader(header, config, baseFileInstantTimeOfPositions), 
keyField);
         blocks.add(dataBlock);
       }
@@ -664,23 +665,8 @@ public class HoodieAppendHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I, K, O
     }
   }
 
-  protected HoodieLogBlock.HoodieLogBlockType pickLogDataBlockFormat() {
-    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 HoodieLogBlock.HoodieLogBlockType.AVRO_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");
-    }
+  protected HoodieLogBlock.HoodieLogBlockType getLogBlockType() {
+    return CommonClientUtils.getLogBlockType(config, 
hoodieTable.getMetaClient().getTableConfig());
   }
 
   private static Map<HeaderMetadataType, String> 
getUpdatedHeader(Map<HeaderMetadataType, String> header,
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
index 19a5f63381b..c52a5129d04 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
@@ -90,6 +90,7 @@ import org.apache.hudi.table.marker.WriteMarkers;
 import org.apache.hudi.table.marker.WriteMarkersFactory;
 import org.apache.hudi.table.storage.HoodieLayoutFactory;
 import org.apache.hudi.table.storage.HoodieStorageLayout;
+import org.apache.hudi.util.CommonClientUtils;
 
 import org.apache.avro.Schema;
 import org.slf4j.Logger;
@@ -964,11 +965,7 @@ public abstract class HoodieTable<T, I, K, O> implements 
Serializable {
   }
 
   public HoodieFileFormat getBaseFileFormat() {
-    HoodieTableConfig tableConfig = metaClient.getTableConfig();
-    if (tableConfig.isMultipleBaseFileFormatsEnabled() && 
config.contains(HoodieWriteConfig.BASE_FILE_FORMAT)) {
-      return config.getBaseFileFormat();
-    }
-    return metaClient.getTableConfig().getBaseFileFormat();
+    return CommonClientUtils.getBaseFileFormat(config, 
metaClient.getTableConfig());
   }
 
   public Option<HoodieFileFormat> getPartitionMetafileFormat() {
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
index 6c30807b8bc..c48a3aa7b0c 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
@@ -22,10 +22,14 @@ package org.apache.hudi.util;
 
 import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.model.HoodieFileFormat;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.table.log.HoodieLogFormat;
+import org.apache.hudi.common.table.log.block.HoodieLogBlock;
+import org.apache.hudi.common.util.Option;
 import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.exception.HoodieNotSupportedException;
 
 import org.slf4j.Logger;
@@ -57,6 +61,37 @@ public class CommonClientUtils {
     }
   }
 
+  /**
+   * Returns the base file format.
+   */
+  public static HoodieFileFormat getBaseFileFormat(HoodieWriteConfig 
writeConfig, HoodieTableConfig tableConfig) {
+    if (tableConfig.isMultipleBaseFileFormatsEnabled() && 
writeConfig.contains(HoodieWriteConfig.BASE_FILE_FORMAT)) {
+      return writeConfig.getBaseFileFormat();
+    }
+    return tableConfig.getBaseFileFormat();
+  }
+
+  /**
+   * Returns the log block type..
+   */
+  public static HoodieLogBlock.HoodieLogBlockType 
getLogBlockType(HoodieWriteConfig writeConfig, HoodieTableConfig tableConfig) {
+    Option<HoodieLogBlock.HoodieLogBlockType> logBlockTypeOpt = 
writeConfig.getLogDataBlockFormat();
+    if (logBlockTypeOpt.isPresent()) {
+      return logBlockTypeOpt.get();
+    }
+    HoodieFileFormat baseFileFormat = getBaseFileFormat(writeConfig, 
tableConfig);
+    switch (getBaseFileFormat(writeConfig, tableConfig)) {
+      case PARQUET:
+      case ORC:
+        return HoodieLogBlock.HoodieLogBlockType.AVRO_DATA_BLOCK;
+      case HFILE:
+        return HoodieLogBlock.HoodieLogBlockType.HFILE_DATA_BLOCK;
+      default:
+        throw new HoodieException("Base file format " + baseFileFormat
+            + " does not have associated log block type");
+    }
+  }
+
   public static String generateWriteToken(TaskContextSupplier 
taskContextSupplier) {
     try {
       return FSUtils.makeWriteToken(
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/HoodieFlinkAvroRecord.java
similarity index 68%
copy from 
hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkRecord.java
copy to 
hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/model/HoodieFlinkAvroRecord.java
index 46ad4a4efd7..a8c9204775d 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/HoodieFlinkAvroRecord.java
@@ -31,42 +31,37 @@ 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.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 org.apache.avro.generic.GenericData;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.IndexedRecord;
 
-import java.io.IOException;
 import java.util.Map;
 import java.util.Properties;
 
 /**
- * Flink Engine-specific Implementations of `HoodieRecord`, which is expected 
to hold {@code RowData} as payload.
+ * 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 HoodieFlinkRecord extends HoodieRecord<RowData> {
+public class HoodieFlinkAvroRecord extends HoodieRecord<IndexedRecord> {
   private Comparable<?> orderingValue = 0;
 
-  public HoodieFlinkRecord(RowData rowData) {
-    super(null, rowData);
-  }
-
-  public HoodieFlinkRecord(HoodieKey key, HoodieOperation op, Comparable<?> 
orderingValue, RowData rowData) {
-    super(key, rowData, op, Option.empty());
+  public HoodieFlinkAvroRecord(HoodieKey key, HoodieOperation op, 
Comparable<?> orderingValue, IndexedRecord record) {
+    super(key, record, op, Option.empty());
     this.orderingValue = orderingValue;
   }
 
   @Override
-  public HoodieRecord<RowData> newInstance() {
-    return new HoodieFlinkRecord(key, operation, orderingValue, data);
+  public HoodieRecord<IndexedRecord> newInstance() {
+    return new HoodieFlinkAvroRecord(key, operation, orderingValue, data);
   }
 
   @Override
-  public HoodieRecord<RowData> newInstance(HoodieKey key, HoodieOperation op) {
-    return new HoodieFlinkRecord(key, op, orderingValue, this.data);
+  public HoodieRecord<IndexedRecord> newInstance(HoodieKey key, 
HoodieOperation op) {
+    return new HoodieFlinkAvroRecord(key, op, orderingValue, data);
   }
 
   @Override
-  public HoodieRecord<RowData> newInstance(HoodieKey key) {
+  public HoodieRecord<IndexedRecord> newInstance(HoodieKey key) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
 
@@ -77,7 +72,7 @@ public class HoodieFlinkRecord extends HoodieRecord<RowData> {
 
   @Override
   public HoodieRecordType getRecordType() {
-    return HoodieRecordType.FLINK;
+    return HoodieRecordType.AVRO;
   }
 
   @Override
@@ -91,12 +86,12 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
   }
 
   @Override
-  protected void writeRecordPayload(RowData payload, Kryo kryo, Output output) 
{
+  protected void writeRecordPayload(IndexedRecord payload, Kryo kryo, Output 
output) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
 
   @Override
-  protected RowData readRecordPayload(Kryo kryo, Input input) {
+  protected IndexedRecord readRecordPayload(Kryo kryo, Input input) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
 
@@ -112,15 +107,22 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
 
   @Override
   public HoodieRecord prependMetaFields(Schema recordSchema, Schema 
targetSchema, MetadataValues metadataValues, Properties props) {
-    int metaFieldSize = targetSchema.getFields().size() - 
recordSchema.getFields().size();
-    GenericRowData metaRow = new GenericRowData(metaFieldSize);
-    String[] metaVals = metadataValues.getValues();
-    for (int i = 0; i < metaVals.length; i++) {
-      if (metaVals[i] != null) {
-        metaRow.setField(i, StringData.fromString(metaVals[i]));
+    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]);
+        }
       }
     }
-    return new HoodieFlinkRecord(key, operation, orderingValue, new 
JoinedRowData(data.getRowKind(), metaRow, data));
+    // 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
@@ -129,7 +131,7 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
   }
 
   @Override
-  public boolean isDelete(Schema recordSchema, Properties props) throws 
IOException {
+  public boolean isDelete(Schema recordSchema, Properties props) {
     if (data == null) {
       return true;
     }
@@ -140,16 +142,20 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
 
     // Use data field to decide.
     Schema.Field deleteField = recordSchema.getField(HOODIE_IS_DELETED_FIELD);
-    return deleteField != null && data.getBoolean(deleteField.pos());
+    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) throws 
IOException {
+  public boolean shouldIgnore(Schema recordSchema, Properties props) {
     return false;
   }
 
   @Override
-  public HoodieRecord<RowData> copy() {
+  public HoodieRecord<IndexedRecord> copy() {
     return this;
   }
 
@@ -160,7 +166,7 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
 
   @Override
   public HoodieRecord wrapIntoHoodieRecordPayloadWithParams(Schema 
recordSchema, Properties props, Option<Pair<String, String>> 
simpleKeyGenFieldsOpt, Boolean withOperation,
-                                                            Option<String> 
partitionNameOp, Boolean populateMetaFieldsOp, Option<Schema> 
schemaWithoutMetaFields) throws IOException {
+                                                            Option<String> 
partitionNameOp, Boolean populateMetaFieldsOp, Option<Schema> 
schemaWithoutMetaFields) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
 
@@ -170,12 +176,12 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
   }
 
   @Override
-  public HoodieRecord truncateRecordKey(Schema recordSchema, Properties props, 
String keyFieldName) throws IOException {
+  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) throws IOException {
-    throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
+  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 46ad4a4efd7..77f0456b17f 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
@@ -129,7 +129,7 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
   }
 
   @Override
-  public boolean isDelete(Schema recordSchema, Properties props) throws 
IOException {
+  public boolean isDelete(Schema recordSchema, Properties props) {
     if (data == null) {
       return true;
     }
@@ -160,7 +160,7 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
 
   @Override
   public HoodieRecord wrapIntoHoodieRecordPayloadWithParams(Schema 
recordSchema, Properties props, Option<Pair<String, String>> 
simpleKeyGenFieldsOpt, Boolean withOperation,
-                                                            Option<String> 
partitionNameOp, Boolean populateMetaFieldsOp, Option<Schema> 
schemaWithoutMetaFields) throws IOException {
+                                                            Option<String> 
partitionNameOp, Boolean populateMetaFieldsOp, Option<Schema> 
schemaWithoutMetaFields) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
 
@@ -170,12 +170,12 @@ public class HoodieFlinkRecord extends 
HoodieRecord<RowData> {
   }
 
   @Override
-  public HoodieRecord truncateRecordKey(Schema recordSchema, Properties props, 
String keyFieldName) throws IOException {
+  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) throws IOException {
+  public Option<HoodieAvroIndexedRecord> toIndexedRecord(Schema recordSchema, 
Properties props) {
     throw new UnsupportedOperationException("Not supported for " + 
this.getClass().getSimpleName());
   }
 }
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 0b517b5d4ae..35b988c142b 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
@@ -24,6 +24,7 @@ import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.action.commit.BucketType;
 import org.apache.hudi.table.marker.WriteMarkers;
 import org.apache.hudi.table.marker.WriteMarkersFactory;
 
@@ -50,6 +51,7 @@ public class FlinkAppendHandle<T, I, K, O>
 
   private boolean isClosed = false;
   private final WriteMarkers writeMarkers;
+  private final BucketType bucketType;
 
   public FlinkAppendHandle(
       HoodieWriteConfig config,
@@ -57,10 +59,12 @@ public class FlinkAppendHandle<T, I, K, O>
       HoodieTable<T, I, K, O> hoodieTable,
       String partitionPath,
       String fileId,
+      BucketType bucketType,
       Iterator<HoodieRecord<T>> recordItr,
       TaskContextSupplier taskContextSupplier) {
     super(config, instantTime, hoodieTable, partitionPath, fileId, recordItr, 
taskContextSupplier);
     this.writeMarkers = WriteMarkersFactory.get(config.getMarkersType(), 
hoodieTable, instantTime);
+    this.bucketType = bucketType;
   }
 
   @Override
@@ -88,8 +92,7 @@ public class FlinkAppendHandle<T, I, K, O>
     // do not use the HoodieRecord operation because hoodie writer has its own
     // INSERT/MERGE bucket for 'UPSERT' semantics. For e.g, a hoodie record 
with fresh new key
     // and operation HoodieCdcOperation.DELETE would be put into either an 
INSERT bucket or UPDATE bucket.
-    return hoodieRecord.getCurrentLocation() != null
-        && hoodieRecord.getCurrentLocation().getInstantTime().equals("U");
+    return bucketType == BucketType.UPDATE;
   }
 
   @Override
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
index 4bc55408cbb..61ba9d88eed 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
@@ -26,6 +26,7 @@ import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.action.commit.BucketType;
 
 import org.apache.hadoop.fs.Path;
 
@@ -280,7 +281,8 @@ public class FlinkWriteHandleFactory {
       final String fileID = record.getCurrentLocation().getFileId();
       final String partitionPath = record.getPartitionPath();
       final TaskContextSupplier contextSupplier = 
table.getTaskContextSupplier();
-      return new FlinkAppendHandle<>(config, instantTime, table, 
partitionPath, fileID, recordItr, contextSupplier);
+      BucketType bucketType = 
record.getCurrentLocation().getInstantTime().equals("I") ? BucketType.INSERT : 
BucketType.UPDATE;
+      return new FlinkAppendHandle<>(config, instantTime, table, 
partitionPath, fileID, bucketType, recordItr, contextSupplier);
     }
   }
 }
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 22788cda57c..27dfd6f6fd1 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,11 +21,14 @@ 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;
 
@@ -66,15 +69,27 @@ public class FlinkRowDataHandleFactory {
         String instantTime,
         HoodieTable<T, I, K, O> table,
         Iterator<HoodieRecord<T>> recordIterator) {
-      return new RowDataLogWriteHandle<>(
-          config,
-          instantTime,
-          table,
-          recordIterator,
-          bucketInfo.getFileIdPrefix(),
-          bucketInfo.getPartitionPath(),
-          bucketInfo.getBucketType(),
-          table.getTaskContextSupplier());
+      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());
+      }
     }
   }
 
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 f89169e75fd..2682d718f30 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
@@ -18,7 +18,6 @@
 
 package org.apache.hudi.io.v2;
 
-import org.apache.hudi.client.WriteStatus;
 import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.model.HoodieColumnRangeMetadata;
@@ -33,14 +32,13 @@ import org.apache.hudi.common.util.SizeEstimator;
 import org.apache.hudi.common.util.ValidationUtils;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieException;
-import org.apache.hudi.io.HoodieAppendHandle;
+import org.apache.hudi.io.FlinkAppendHandle;
 import org.apache.hudi.io.MiniBatchHandle;
 import org.apache.hudi.io.log.block.HoodieFlinkParquetDataBlock;
 import org.apache.hudi.io.storage.ColumnRangeMetadataProvider;
 import org.apache.hudi.io.storage.row.HoodieFlinkIOFactory;
 import org.apache.hudi.metadata.HoodieTableMetadataUtil;
 import org.apache.hudi.storage.StorageConfiguration;
-import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
 import org.apache.hudi.table.action.commit.BucketType;
 import org.apache.hudi.util.Lazy;
@@ -71,14 +69,10 @@ import static 
org.apache.hudi.common.config.HoodieStorageConfig.PARQUET_DICTIONA
  * given file group and instant.
  */
 public class RowDataLogWriteHandle<T, I, K, O>
-    extends HoodieAppendHandle<T, I, K, O> implements MiniBatchHandle {
+    extends FlinkAppendHandle<T, I, K, O> implements MiniBatchHandle {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(RowDataLogWriteHandle.class);
 
-  private boolean isClosed = false;
-
-  private final BucketType bucketType;
-
   public RowDataLogWriteHandle(
       HoodieWriteConfig config,
       String instantTime,
@@ -88,9 +82,8 @@ public class RowDataLogWriteHandle<T, I, K, O>
       String partitionPath,
       BucketType bucketType,
       TaskContextSupplier taskContextSupplier) {
-    super(config, instantTime, hoodieTable, partitionPath, fileId, recordItr, 
taskContextSupplier);
+    super(config, instantTime, hoodieTable, partitionPath, fileId, bucketType, 
recordItr, taskContextSupplier);
     initWriteConf(storage.getConf(), config);
-    this.bucketType = bucketType;
   }
 
   private void initWriteConf(StorageConfiguration<?> storageConf, 
HoodieWriteConfig writeConfig) {
@@ -107,6 +100,25 @@ 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.
    */
@@ -186,54 +198,4 @@ public class RowDataLogWriteHandle<T, I, K, O>
         throw new HoodieException("Data block format " + logDataBlockFormat + 
" is not implemented for Flink RowData append handle.");
     }
   }
-
-  @Override
-  protected boolean isUpdateRecord(HoodieRecord<T> hoodieRecord) {
-    return bucketType == BucketType.UPDATE;
-  }
-
-  @Override
-  protected boolean needsUpdateLocation() {
-    return false;
-  }
-
-  @Override
-  public boolean canWrite(HoodieRecord record) {
-    return true;
-  }
-
-  @Override
-  protected HoodieLogBlock.HoodieLogBlockType pickLogDataBlockFormat() {
-    Option<HoodieLogBlock.HoodieLogBlockType> logBlockTypeOpt = 
config.getLogDataBlockFormat();
-    if (logBlockTypeOpt.isPresent()) {
-      return logBlockTypeOpt.get();
-    }
-    return HoodieLogBlock.HoodieLogBlockType.PARQUET_DATA_BLOCK;
-  }
-
-  @Override
-  public List<WriteStatus> close() {
-    try {
-      return super.close();
-    } finally {
-      this.isClosed = true;
-    }
-  }
-
-  @Override
-  public void closeGracefully() {
-    if (isClosed) {
-      return;
-    }
-    try {
-      close();
-    } catch (Throwable throwable) {
-      LOG.warn("Error while trying to dispose the APPEND handle", throwable);
-    }
-  }
-
-  @Override
-  public StoragePath getWritePath() {
-    return writer.getLogFile().getPath();
-  }
 }
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
index 54d8225a24b..c0a2579e248 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
@@ -34,7 +34,6 @@ import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.io.FlinkAppendHandle;
 import org.apache.hudi.io.HoodieAppendHandle;
 import org.apache.hudi.io.HoodieWriteHandle;
-import org.apache.hudi.io.v2.RowDataLogWriteHandle;
 import org.apache.hudi.table.action.HoodieWriteMetadata;
 import org.apache.hudi.table.action.commit.BucketInfo;
 import 
org.apache.hudi.table.action.commit.delta.FlinkUpsertDeltaCommitActionExecutor;
@@ -83,10 +82,9 @@ public class HoodieFlinkMergeOnReadTable<T>
       BucketInfo bucketInfo,
       String instantTime,
       Iterator<HoodieRecord<T>> records) {
-    ValidationUtils.checkArgument(writeHandle instanceof RowDataLogWriteHandle,
-        "MOR RowData handle should always be a RowDataLogHandle");
-    RowDataLogWriteHandle<?, ?, ?, ?> rowDataLogHandle = 
(RowDataLogWriteHandle<?, ?, ?, ?>) writeHandle;
-    return new RowDataUpsertDeltaCommitActionExecutor<>(context, 
rowDataLogHandle, bucketInfo, config, this, instantTime, records).execute();
+    ValidationUtils.checkArgument(writeHandle instanceof FlinkAppendHandle,
+        "MOR RowData handle should always be a FlinkAppendHandle");
+    return new RowDataUpsertDeltaCommitActionExecutor<>(context, writeHandle, 
bucketInfo, config, this, instantTime, records).execute();
   }
 
   @Override
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/delta/RowDataUpsertDeltaCommitActionExecutor.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/delta/RowDataUpsertDeltaCommitActionExecutor.java
index bad62bef0db..dbb1a02f729 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/delta/RowDataUpsertDeltaCommitActionExecutor.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/delta/RowDataUpsertDeltaCommitActionExecutor.java
@@ -23,7 +23,8 @@ import org.apache.hudi.common.engine.HoodieEngineContext;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.config.HoodieWriteConfig;
-import org.apache.hudi.io.v2.RowDataLogWriteHandle;
+import org.apache.hudi.io.FlinkAppendHandle;
+import org.apache.hudi.io.HoodieWriteHandle;
 import org.apache.hudi.table.HoodieTable;
 import org.apache.hudi.table.action.HoodieWriteMetadata;
 import org.apache.hudi.table.action.commit.BaseFlinkCommitActionExecutor;
@@ -43,7 +44,7 @@ public class RowDataUpsertDeltaCommitActionExecutor<T> 
extends BaseFlinkCommitAc
   private final BucketInfo bucketInfo;
 
   public RowDataUpsertDeltaCommitActionExecutor(HoodieEngineContext context,
-                                              RowDataLogWriteHandle<?, ?, ?, 
?> writeHandle,
+                                              HoodieWriteHandle<?, ?, ?, ?> 
writeHandle,
                                               BucketInfo bucketInfo,
                                               HoodieWriteConfig config,
                                               HoodieTable table,
@@ -68,7 +69,7 @@ public class RowDataUpsertDeltaCommitActionExecutor<T> 
extends BaseFlinkCommitAc
    * When using RowData writing mode for MOR upsert path, same write logic is 
used for UPSERT and INSERT operation.
    */
   private Iterator<List<WriteStatus>> handleWrite() {
-    RowDataLogWriteHandle logWriteHandle = (RowDataLogWriteHandle) writeHandle;
+    FlinkAppendHandle logWriteHandle = (FlinkAppendHandle) writeHandle;
     logWriteHandle.doAppend();
     List<WriteStatus> writeStatuses = logWriteHandle.close();
     return Collections.singletonList(writeStatuses).iterator();
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/OrderingValueExtractor.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/OrderingValueExtractor.java
deleted file mode 100644
index e404fcba8ca..00000000000
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/util/OrderingValueExtractor.java
+++ /dev/null
@@ -1,30 +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.flink.table.data.RowData;
-
-import java.io.Serializable;
-
-/**
- * Interface for extracting ordering value from RowData.
- */
-public interface OrderingValueExtractor extends Serializable {
-  Comparable<?> getOrderingValue(RowData rowData);
-}
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 0a440919ab6..f7a47bd4098 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
@@ -18,6 +18,7 @@
 
 package org.apache.hudi.common.table.log.block;
 
+import org.apache.hudi.avro.AvroSchemaCache;
 import org.apache.hudi.avro.HoodieAvroUtils;
 import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.fs.SizeAwareDataInputStream;
@@ -91,8 +92,7 @@ public class HoodieAvroDataBlock extends HoodieDataBlock {
 
   public HoodieAvroDataBlock(@Nonnull List<HoodieRecord> records,
                              @Nonnull Map<HeaderMetadataType, String> header,
-                             @Nonnull String keyField
-  ) {
+                             @Nonnull String keyField) {
     super(records, header, new HashMap<>(), keyField);
   }
 
@@ -103,7 +103,7 @@ public class HoodieAvroDataBlock extends HoodieDataBlock {
 
   @Override
   protected byte[] serializeRecords(List<HoodieRecord> records, HoodieStorage 
storage) throws IOException {
-    Schema schema = new 
Schema.Parser().parse(super.getLogBlockHeader().get(HeaderMetadataType.SCHEMA));
+    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)) {
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 cbe492b92b4..4f69dfe91c0 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
@@ -18,22 +18,17 @@
 
 package org.apache.hudi.sink;
 
-import org.apache.hudi.avro.HoodieAvroUtils;
 import org.apache.hudi.client.WriteStatus;
 import org.apache.hudi.client.model.HoodieFlinkInternalRow;
-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.model.HoodieRecordMerger;
 import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.common.util.HoodieRecordUtils;
-import org.apache.hudi.common.util.StringUtils;
 import org.apache.hudi.common.util.ValidationUtils;
 import org.apache.hudi.common.util.VisibleForTesting;
 import org.apache.hudi.common.util.collection.MappingIterator;
 import org.apache.hudi.configuration.FlinkOptions;
-import org.apache.hudi.configuration.OptionsResolver;
 import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.metrics.FlinkStreamWriteMetrics;
 import org.apache.hudi.sink.buffer.MemorySegmentPoolFactory;
@@ -43,23 +38,18 @@ import org.apache.hudi.sink.bulk.RowDataKeyGen;
 import org.apache.hudi.sink.common.AbstractStreamWriteFunction;
 import org.apache.hudi.sink.event.WriteMetadataEvent;
 import org.apache.hudi.sink.exception.MemoryPagesExhaustedException;
+import org.apache.hudi.sink.transform.RecordConverter;
 import org.apache.hudi.sink.utils.BufferUtils;
 import org.apache.hudi.table.action.commit.BucketInfo;
 import org.apache.hudi.table.action.commit.BucketType;
-import org.apache.hudi.util.AvroSchemaConverter;
 import org.apache.hudi.util.MutableIteratorWrapperIterator;
-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.flink.configuration.Configuration;
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.streaming.api.functions.ProcessFunction;
-import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.data.binary.BinaryRowData;
 import org.apache.flink.table.runtime.util.MemorySegmentPool;
-import org.apache.flink.table.types.logical.LogicalType;
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.types.RowKind;
 import org.apache.flink.util.Collector;
@@ -67,6 +57,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.io.IOException;
+import java.io.Serializable;
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.Comparator;
@@ -130,8 +121,6 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
 
   protected final RowDataKeyGen keyGen;
 
-  protected final OrderingValueExtractor orderingValueExtractor;
-
   /**
    * Total size tracer.
    */
@@ -144,6 +133,8 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
 
   protected transient MemorySegmentPool memorySegmentPool;
 
+  protected transient RecordConverter recordConverter;
+
   /**
    * Constructs a StreamingSinkFunction.
    *
@@ -153,7 +144,6 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
     super(config);
     this.rowType = rowType;
     this.keyGen = RowDataKeyGen.instance(config, rowType);
-    this.orderingValueExtractor = getOrderingValueExtractor(config, rowType);
   }
 
   @Override
@@ -162,6 +152,7 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
     initBuffer();
     initWriteFunction();
     initMergeClass();
+    initRecordConverter();
     registerMetrics();
   }
 
@@ -227,6 +218,10 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
     }
   }
 
+  private void initRecordConverter() {
+    this.recordConverter = RecordConverter.getInstance(config, rowType, 
keyGen, writeClient.getConfig(), metaClient.getTableConfig());
+  }
+
   private void initMergeClass() {
     recordMerger = 
HoodieRecordUtils.mergerToPreCombineMode(writeClient.getConfig().getRecordMerger());
     LOG.info("init hoodie merge with class [{}]", 
recordMerger.getClass().getName());
@@ -279,19 +274,21 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
         
RowKind.fromByteValue(HoodieOperation.fromName(record.getOperationType()).getValue()));
     final String bucketID = getBucketID(record.getPartitionPath(), 
record.getFileId());
 
+    // 1. try buffer the record into the memory pool
     boolean success = doBufferRecord(bucketID, record);
-    // 1. flushing bucket for memory pool is full.
     if (!success) {
+      // 2. flushes the bucket if the memory pool is full
       RowDataBucket bucketToFlush = this.buckets.values().stream()
           .max(Comparator.comparingLong(RowDataBucket::getBufferSize))
           .orElseThrow(NoSuchElementException::new);
       if (flushBucket(bucketToFlush)) {
+        // 2.1 flushes the data bucket with maximum size
         this.tracer.countDown(bucketToFlush.getBufferSize());
         disposeBucket(bucketToFlush);
       } else {
         LOG.warn("The buffer size hits the threshold {}, but still flush the 
max size data bucket failed!", this.tracer.maxBufferSize);
       }
-      // try write row again
+      // 2.2 try to write row again
       success = doBufferRecord(bucketID, record);
       if (!success) {
         throw new RuntimeException("Buffer is too small to hold a single 
record.");
@@ -299,7 +296,7 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
     }
     RowDataBucket bucket = this.buckets.get(bucketID);
     this.tracer.trace(bucket.getLastRecordSize());
-    // 2. flushing bucket for bucket is full.
+    // 3. flushes the bucket if it is full
     if (bucket.isFull()) {
       if (flushBucket(bucket)) {
         this.tracer.countDown(bucket.getBufferSize());
@@ -405,12 +402,6 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
     writeMetrics.resetAfterCommit();
   }
 
-  private void registerMetrics() {
-    MetricGroup metrics = getRuntimeContext().getMetricGroup();
-    writeMetrics = new FlinkStreamWriteMetrics(metrics);
-    writeMetrics.registerMetrics();
-  }
-
   protected List<WriteStatus> writeRecords(
       String instant,
       RowDataBucket rowDataBucket) {
@@ -420,7 +411,7 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
         new MutableIteratorWrapperIterator<>(
             rowDataBucket.getDataIterator(), () -> new 
BinaryRowData(rowType.getFieldCount()));
     Iterator<HoodieRecord> recordItr = new MappingIterator<>(
-        rowItr, rowData -> convertToRecord(rowData, 
rowDataBucket.getBucketInfo()));
+        rowItr, rowData -> recordConverter.convert(rowData, 
rowDataBucket.getBucketInfo()));
 
     List<WriteStatus> statuses = 
writeFunction.write(deduplicateRecordsIfNeeded(recordItr), 
rowDataBucket.getBucketInfo(), instant);
     writeMetrics.endFileFlush();
@@ -428,14 +419,6 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
     return statuses;
   }
 
-  protected HoodieFlinkRecord convertToRecord(RowData dataRow, BucketInfo 
bucketInfo) {
-    String key = keyGen.getRecordKey(dataRow);
-    Comparable<?> preCombineValue = 
orderingValueExtractor.getOrderingValue(dataRow);
-    HoodieOperation operation = 
HoodieOperation.fromValue(dataRow.getRowKind().toByteValue());
-    HoodieKey hoodieKey = new HoodieKey(key, bucketInfo.getPartitionPath());
-    return new HoodieFlinkRecord(hoodieKey, operation, preCombineValue, 
dataRow);
-  }
-
   protected Iterator<HoodieRecord> 
deduplicateRecordsIfNeeded(Iterator<HoodieRecord> records) {
     if (config.get(FlinkOptions.PRE_COMBINE)) {
       // todo: sort by record key, lazy merge during iterating, default for 
COW.
@@ -445,34 +428,10 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
     }
   }
 
-  private OrderingValueExtractor getOrderingValueExtractor(Configuration conf, 
RowType rowType) {
-    String preCombineField = OptionsResolver.getPreCombineField(conf);
-    if (StringUtils.isNullOrEmpty(preCombineField)) {
-      // return a dummy extractor.
-      return new OrderingValueExtractor() {
-        @Override
-        public Comparable<?> getOrderingValue(RowData rowData) {
-          return HoodieRecord.DEFAULT_ORDERING_VALUE;
-        }
-      };
-    }
-    int preCombineFieldIdx = rowType.getFieldNames().indexOf(preCombineField);
-    LogicalType fieldType = rowType.getChildren().get(preCombineFieldIdx);
-    RowData.FieldGetter preCombineFieldGetter = 
RowData.createFieldGetter(fieldType, preCombineFieldIdx);
-
-    // currently the log reader for flink is avro reader, and it merges 
records based on ordering value
-    // in form of AVRO format, so here we align the data format with reader.
-    // TODO refactor this after RowData reader is supported, HUDI-9146.
-    RowDataToAvroConverters.RowDataToAvroConverter fieldConverter =
-        RowDataToAvroConverters.createConverter(fieldType, 
conf.get(FlinkOptions.WRITE_UTC_TIMEZONE));
-    Schema fieldSchema = AvroSchemaConverter.convertToSchema(fieldType, 
preCombineField);
-    return new OrderingValueExtractor() {
-      @Override
-      public Comparable<?> getOrderingValue(RowData rowData) {
-        return (Comparable<?>) 
HoodieAvroUtils.convertValueForSpecificDataTypes(
-            fieldSchema, fieldConverter.convert(fieldSchema, 
preCombineFieldGetter.getFieldOrNull(rowData)), false);
-      }
-    };
+  private void registerMetrics() {
+    MetricGroup metrics = getRuntimeContext().getMetricGroup();
+    writeMetrics = new FlinkStreamWriteMetrics(metrics);
+    writeMetrics.registerMetrics();
   }
 
   // -------------------------------------------------------------------------
@@ -489,14 +448,21 @@ public class RowDataStreamWriteFunction extends 
AbstractStreamWriteFunction<Hood
           new MutableIteratorWrapperIterator<>(
               entry.getValue().getDataIterator(), () -> new 
BinaryRowData(rowType.getFieldCount()));
       while (rowItr.hasNext()) {
-        records.add(convertToRecord(rowItr.next(), 
entry.getValue().getBucketInfo()));
+        records.add(recordConverter.convert(rowItr.next(), 
entry.getValue().getBucketInfo()));
       }
       ret.put(entry.getKey(), records);
     }
     return ret;
   }
 
-  protected interface WriteFunction {
+  // -------------------------------------------------------------------------
+  //  Inner Classes
+  // -------------------------------------------------------------------------
+
+  /**
+   * Write function to trigger the actual write action.
+   */
+  protected interface WriteFunction extends Serializable {
     List<WriteStatus> write(Iterator<HoodieRecord> records, BucketInfo 
bucketInfo, String instant);
   }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bucket/RowDataConsistentBucketStreamWriteFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bucket/RowDataConsistentBucketStreamWriteFunction.java
index 40ce155f389..a5f328cd3af 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bucket/RowDataConsistentBucketStreamWriteFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bucket/RowDataConsistentBucketStreamWriteFunction.java
@@ -87,7 +87,7 @@ public class RowDataConsistentBucketStreamWriteFunction 
extends RowDataStreamWri
         new MutableIteratorWrapperIterator<>(
             rowDataBucket.getDataIterator(), () -> new 
BinaryRowData(rowType.getFieldCount()));
     Iterator<HoodieRecord> recordItr = deduplicateRecordsIfNeeded(
-        new MappingIterator<>(rowItr, rowData -> convertToRecord(rowData, 
rowDataBucket.getBucketInfo())));
+        new MappingIterator<>(rowItr, rowData -> 
recordConverter.convert(rowData, rowDataBucket.getBucketInfo())));
 
     Pair<List<BucketRecords>, Set<HoodieFileGroupId>> recordListFgPair =
         
updateStrategy.handleUpdate(Collections.singletonList(BucketRecords.of(recordItr,
 rowDataBucket.getBucketInfo(), instant)));
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
new file mode 100644
index 00000000000..86095730376
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/transform/RecordConverter.java
@@ -0,0 +1,89 @@
+/*
+ * 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.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;
+
+/**
+ * Function that converts the given {@link RowData} into a hoodie record.
+ */
+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);
+    }
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/DataTypeUtils.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/DataTypeUtils.java
index bddef13eeeb..b2b8a1cc7f6 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/DataTypeUtils.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/DataTypeUtils.java
@@ -19,8 +19,12 @@
 package org.apache.hudi.util;
 
 import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.exception.HoodieValidationException;
 
+import org.apache.avro.util.Utf8;
 import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.data.DecimalData;
+import org.apache.flink.table.data.TimestampData;
 import org.apache.flink.table.types.DataType;
 import org.apache.flink.table.types.logical.LocalZonedTimestampType;
 import org.apache.flink.table.types.logical.LogicalType;
@@ -32,6 +36,8 @@ import org.apache.flink.table.types.logical.TimestampType;
 import javax.annotation.Nullable;
 
 import java.math.BigDecimal;
+import java.nio.ByteBuffer;
+import java.time.Instant;
 import java.time.LocalDate;
 import java.time.LocalDateTime;
 import java.util.ArrayList;
@@ -205,4 +211,63 @@ public class DataTypeUtils {
 
     return new RowType(false, mergedFields);
   }
+
+  /**
+   * 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 fieldVal    The field value
+   * @param utcTimezone whether to use UTC timezone for timestamp data type
+   */
+  public static Object resolveOrderingValue(
+      LogicalType logicalType,
+      Object fieldVal,
+      boolean utcTimezone) {
+    switch (logicalType.getTypeRoot()) {
+      case NULL:
+        return null;
+      case TINYINT:
+        return ((Byte) fieldVal).intValue();
+      case SMALLINT:
+        return ((Short) fieldVal).intValue();
+      case DATE:
+        return LocalDate.ofEpochDay((Long) fieldVal);
+      case CHAR:
+      case VARCHAR:
+        return new Utf8(fieldVal.toString());
+      case BINARY:
+      case VARBINARY:
+        return ByteBuffer.wrap((byte[]) fieldVal);
+      case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+        int precision1 = DataTypeUtils.precision(logicalType);
+        if (precision1 <= 3) {
+          return ((TimestampData) fieldVal).toInstant().toEpochMilli();
+        } else if (precision1 <= 6) {
+          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 = DataTypeUtils.precision(logicalType);
+        if (precision2 <= 3) {
+          return utcTimezone ? ((TimestampData) 
fieldVal).toInstant().toEpochMilli() : ((TimestampData) 
fieldVal).toTimestamp().getTime();
+        } else if (precision2 <= 6) {
+          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 ((DecimalData) fieldVal).toBigDecimal();
+      default:
+        if (fieldVal == null) {
+          throw new HoodieValidationException("Ordering value(legacy as 
preCombine field value) can not be null");
+        }
+        return fieldVal;
+    }
+  }
 }
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
new file mode 100644
index 00000000000..fdebb738a9b
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/OrderingValueExtractor.java
@@ -0,0 +1,57 @@
+/*
+ * 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 560987711ad..2d0478a484a 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
@@ -19,6 +19,7 @@
 package org.apache.hudi.table;
 
 import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
 import org.apache.hudi.common.model.HoodieTableType;
 import org.apache.hudi.common.model.WriteOperationType;
@@ -1926,6 +1927,46 @@ public class ITTestHoodieDataSource {
         + "+I[id8, Han, 56, 1970-01-01T00:00:08, par4]]");
   }
 
+  @Test
+  void testParquetLogBlockDataSkipping() {
+    TableEnvironment tableEnv = batchTableEnv;
+    String hoodieTableDDL = sql("t1")
+        .option(FlinkOptions.PATH, tempFile.getAbsolutePath())
+        .option(FlinkOptions.METADATA_ENABLED, true)
+        .option("hoodie.metadata.index.column.stats.enable", true)
+        .option(HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key(), "parquet")
+        .option(FlinkOptions.READ_DATA_SKIPPING_ENABLED, true)
+        .option(FlinkOptions.TABLE_TYPE, FlinkOptions.TABLE_TYPE_MERGE_ON_READ)
+        .end();
+    tableEnv.executeSql(hoodieTableDDL);
+
+    execInsertSql(tableEnv, TestSQL.INSERT_T1);
+
+    List<Row> result1 = CollectionUtil.iterableToList(
+        () -> tableEnv.sqlQuery("select * from t1").execute().collect());
+    assertRowsEquals(result1, TestData.DATA_SET_SOURCE_INSERT);
+    // apply filters
+    List<Row> result2 = CollectionUtil.iterableToList(
+        () -> tableEnv.sqlQuery("select * from t1 where uuid > 'id5' and age > 
20").execute().collect());
+    assertRowsEquals(result2, "["
+        + "+I[id7, Bob, 44, 1970-01-01T00:00:07, par4], "
+        + "+I[id8, Han, 56, 1970-01-01T00:00:08, par4]]");
+    // filter by timestamp
+    List<Row> result3 = CollectionUtil.iterableToList(
+        () -> tableEnv.sqlQuery("select * from t1 where ts > TIMESTAMP 
'1970-01-01 00:00:05'").execute().collect());
+    assertRowsEquals(result3, "["
+        + "+I[id6, Emma, 20, 1970-01-01T00:00:06, par3], "
+        + "+I[id7, Bob, 44, 1970-01-01T00:00:07, par4], "
+        + "+I[id8, Han, 56, 1970-01-01T00:00:08, par4]]");
+    // filter by in expression
+    List<Row> result4 = CollectionUtil.iterableToList(
+        () -> tableEnv.sqlQuery("select * from t1 where uuid in ('id6', 'id7', 
'id8')").execute().collect());
+    assertRowsEquals(result4, "["
+        + "+I[id6, Emma, 20, 1970-01-01T00:00:06, par3], "
+        + "+I[id7, Bob, 44, 1970-01-01T00:00:07, par4], "
+        + "+I[id8, Han, 56, 1970-01-01T00:00:08, par4]]");
+  }
+
   @Disabled("for being flaky by HUDI-7174")
   @Test
   void testMultipleLogBlocksWithDataSkipping() {
@@ -2508,6 +2549,29 @@ public class ITTestHoodieDataSource {
     assertRowsEquals(rows, TestData.DATA_SET_SOURCE_INSERT);
   }
 
+  @ParameterizedTest
+  @ValueSource(strings = {"FLINK_STATE", "BUCKET"})
+  void testRowDataWriteModeWithParquetLogFormat(String index) throws Exception 
{
+    String createSource = TestConfigurations.getFileSourceDDL("source");
+    streamTableEnv.executeSql(createSource);
+
+    // insert first batch of data with rowdata mode writing disabled
+    String hoodieTableDDL = sql("t1")
+        .option(FlinkOptions.PATH, tempFile.getAbsolutePath())
+        .option(FlinkOptions.TABLE_TYPE, HoodieTableType.MERGE_ON_READ)
+        .option(FlinkOptions.INDEX_TYPE, index)
+        .option(HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key(), "parquet")
+        .option(HoodieWriteConfig.ALLOW_EMPTY_COMMIT.key(), false)
+        .end();
+    streamTableEnv.executeSql(hoodieTableDDL);
+    String insertInto = "insert into t1 select * from source";
+    execInsertSql(streamTableEnv, insertInto);
+
+    // reading from the earliest
+    List<Row> rows = execSelectSqlWithExpectedNum(streamTableEnv, "select * 
from t1", TestData.DATA_SET_SOURCE_INSERT.size());
+    assertRowsEquals(rows, TestData.DATA_SET_SOURCE_INSERT);
+  }
+
   // -------------------------------------------------------------------------
   //  Utilities
   // -------------------------------------------------------------------------
@@ -2554,6 +2618,19 @@ 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()),

Reply via email to