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 f47ef8067fe5 feat: support read native cdc logs (#19114)
f47ef8067fe5 is described below

commit f47ef8067fe51fffca0377f050519baaab9c5d93
Author: Danny Chan <[email protected]>
AuthorDate: Wed Jul 1 13:43:31 2026 +0800

    feat: support read native cdc logs (#19114)
---
 .../table/log/HoodieCDCEngineRecordAccessor.java   |  33 ++++
 ....java => HoodieCDCInlineLogRecordIterator.java} |  70 ++++++--
 .../hudi/common/table/log/HoodieCDCLogRecord.java  |  40 +++++
 .../table/log/HoodieCDCLogRecordIterator.java      | 120 ++-----------
 .../log/HoodieCDCNativeLogRecordIterator.java      | 117 ++++++++++++
 .../log/TestHoodieCDCNativeLogRecordIterator.java  | 162 +++++++++++++++++
 .../function/HoodieCdcSplitReaderFunction.java     |   6 +-
 .../hudi/table/format/cdc/CdcInputFormat.java      |   6 +-
 .../apache/hudi/table/format/cdc/CdcIterators.java | 200 +++++++++++++++++----
 .../org/apache/hudi/cdc/CDCFileGroupIterator.scala | 181 ++++++++++++-------
 10 files changed, 706 insertions(+), 229 deletions(-)

diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCEngineRecordAccessor.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCEngineRecordAccessor.java
new file mode 100644
index 000000000000..2502bb2ea810
--- /dev/null
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCEngineRecordAccessor.java
@@ -0,0 +1,33 @@
+/*
+ * 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.common.table.log;
+
+/**
+ * Accessor for CDC operation, record key, and before/after images in an 
engine-specific row.
+ *
+ * @param <T> Engine-specific record type used by native CDC log files
+ */
+public interface HoodieCDCEngineRecordAccessor<T> {
+
+  String getOperation(T record);
+
+  String getRecordKey(T record);
+
+  T getImage(T record, int ordinal, int imageArity);
+}
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecordIterator.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCInlineLogRecordIterator.java
similarity index 59%
copy from 
hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecordIterator.java
copy to 
hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCInlineLogRecordIterator.java
index b1ccb6019465..bc947e998a8e 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecordIterator.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCInlineLogRecordIterator.java
@@ -33,11 +33,12 @@ import org.apache.avro.generic.IndexedRecord;
 import java.io.IOException;
 import java.util.Arrays;
 import java.util.Iterator;
+import java.util.NoSuchElementException;
 
 /**
- * Record iterator for Hudi logs in CDC format.
+ * CDC log record iterator for inline CDC log blocks.
  */
-public class HoodieCDCLogRecordIterator implements 
ClosableIterator<IndexedRecord> {
+public class HoodieCDCInlineLogRecordIterator implements 
HoodieCDCLogRecordIterator<IndexedRecord> {
 
   private final HoodieStorage storage;
 
@@ -47,11 +48,11 @@ public class HoodieCDCLogRecordIterator implements 
ClosableIterator<IndexedRecor
 
   private HoodieLogFormat.Reader reader;
 
-  private ClosableIterator<IndexedRecord> itr;
+  private ClosableIterator<HoodieCDCLogRecord<IndexedRecord>> itr;
 
-  private IndexedRecord record;
+  private HoodieCDCLogRecord<IndexedRecord> record;
 
-  public HoodieCDCLogRecordIterator(HoodieStorage storage, HoodieLogFile[] 
cdcLogFiles, HoodieSchema cdcSchema) {
+  public HoodieCDCInlineLogRecordIterator(HoodieStorage storage, 
HoodieLogFile[] cdcLogFiles, HoodieSchema cdcSchema) {
     this.storage = storage;
     this.cdcSchema = cdcSchema;
     this.cdcLogFileIter = Arrays.stream(cdcLogFiles).iterator();
@@ -64,12 +65,10 @@ public class HoodieCDCLogRecordIterator implements 
ClosableIterator<IndexedRecor
     }
     if (itr == null || !itr.hasNext()) {
       if (reader == null || !reader.hasNext()) {
-        // step1: load new file reader first.
         if (!loadReader()) {
           return false;
         }
       }
-      // step2: load block records iterator
       if (!loadItr()) {
         return false;
       }
@@ -82,26 +81,32 @@ public class HoodieCDCLogRecordIterator implements 
ClosableIterator<IndexedRecor
     try {
       closeReader();
       if (cdcLogFileIter.hasNext()) {
-        reader = new HoodieLogFileReader(storage, cdcLogFileIter.next(), 
cdcSchema, HoodieLogFileReader.DEFAULT_BUFFER_SIZE);
+        reader = new HoodieLogFileReader(
+            storage, cdcLogFileIter.next(), cdcSchema, 
HoodieLogFileReader.DEFAULT_BUFFER_SIZE);
         return reader.hasNext();
       }
       return false;
     } catch (IOException e) {
-      throw new HoodieIOException(e.getMessage());
+      throw new HoodieIOException(e.getMessage(), e);
     }
   }
 
   private boolean loadItr() {
     HoodieDataBlock dataBlock = (HoodieDataBlock) reader.next();
     closeItr();
-    // TODO support cdc with spark record.
-    itr = new 
CloseableMappingIterator(dataBlock.getRecordIterator(HoodieRecordType.AVRO), 
record -> ((HoodieAvroIndexedRecord) record).getData());
+    itr = new CloseableMappingIterator<>(
+        dataBlock.getRecordIterator(HoodieRecordType.AVRO),
+        // Cast via Object to avoid an unchecked-cast warning; AVRO record 
iterators yield HoodieAvroIndexedRecord.
+        record -> new InlineCDCLogRecord(((HoodieAvroIndexedRecord) (Object) 
record).getData()));
     return itr.hasNext();
   }
 
   @Override
-  public IndexedRecord next() {
-    IndexedRecord ret = record;
+  public HoodieCDCLogRecord<IndexedRecord> next() {
+    if (!hasNext()) {
+      throw new NoSuchElementException("No more CDC log records");
+    }
+    HoodieCDCLogRecord<IndexedRecord> ret = record;
     record = null;
     return ret;
   }
@@ -112,14 +117,10 @@ public class HoodieCDCLogRecordIterator implements 
ClosableIterator<IndexedRecor
       closeItr();
       closeReader();
     } catch (IOException e) {
-      throw new HoodieIOException(e.getMessage());
+      throw new HoodieIOException(e.getMessage(), e);
     }
   }
 
-  // -------------------------------------------------------------------------
-  //  Utilities
-  // -------------------------------------------------------------------------
-
   private void closeReader() throws IOException {
     if (reader != null) {
       reader.close();
@@ -133,4 +134,37 @@ public class HoodieCDCLogRecordIterator implements 
ClosableIterator<IndexedRecor
       itr = null;
     }
   }
+
+  private static class InlineCDCLogRecord implements 
HoodieCDCLogRecord<IndexedRecord> {
+    private final IndexedRecord record;
+
+    private InlineCDCLogRecord(IndexedRecord record) {
+      this.record = record;
+    }
+
+    @Override
+    public String getOperation() {
+      return String.valueOf(record.get(0));
+    }
+
+    @Override
+    public String getRecordKey() {
+      return String.valueOf(record.get(1));
+    }
+
+    @Override
+    public IndexedRecord getAvroImage(int ordinal) {
+      return (IndexedRecord) record.get(ordinal);
+    }
+
+    @Override
+    public IndexedRecord getEngineImage(int ordinal, int imageArity) {
+      throw new UnsupportedOperationException("Inline CDC records do not 
contain engine row images");
+    }
+
+    @Override
+    public boolean isNative() {
+      return false;
+    }
+  }
 }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecord.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecord.java
new file mode 100644
index 000000000000..40c9113410d3
--- /dev/null
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecord.java
@@ -0,0 +1,40 @@
+/*
+ * 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.common.table.log;
+
+import org.apache.avro.generic.IndexedRecord;
+
+/**
+ * Normalized CDC log record that can represent inline Avro CDC records or 
native CDC records
+ * materialized in an engine-specific row type.
+ *
+ * @param <T> Engine-specific record type used by native CDC log files
+ */
+public interface HoodieCDCLogRecord<T> {
+
+  String getOperation();
+
+  String getRecordKey();
+
+  IndexedRecord getAvroImage(int ordinal);
+
+  T getEngineImage(int ordinal, int imageArity);
+
+  boolean isNative();
+}
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecordIterator.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecordIterator.java
index b1ccb6019465..169f03a2ed4c 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecordIterator.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCLogRecordIterator.java
@@ -18,119 +18,33 @@
 
 package org.apache.hudi.common.table.log;
 
-import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
-import org.apache.hudi.common.model.HoodieLogFile;
-import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType;
-import org.apache.hudi.common.schema.HoodieSchema;
-import org.apache.hudi.common.table.log.block.HoodieDataBlock;
 import org.apache.hudi.common.util.collection.ClosableIterator;
-import org.apache.hudi.common.util.collection.CloseableMappingIterator;
-import org.apache.hudi.exception.HoodieIOException;
-import org.apache.hudi.storage.HoodieStorage;
 
-import org.apache.avro.generic.IndexedRecord;
-
-import java.io.IOException;
-import java.util.Arrays;
-import java.util.Iterator;
+import java.util.NoSuchElementException;
 
 /**
- * Record iterator for Hudi logs in CDC format.
+ * Record iterator for Hudi CDC log files.
+ *
+ * @param <T> Engine-specific record type used by native CDC log files
  */
-public class HoodieCDCLogRecordIterator implements 
ClosableIterator<IndexedRecord> {
-
-  private final HoodieStorage storage;
-
-  private final HoodieSchema cdcSchema;
-
-  private final Iterator<HoodieLogFile> cdcLogFileIter;
-
-  private HoodieLogFormat.Reader reader;
+public interface HoodieCDCLogRecordIterator<T> extends 
ClosableIterator<HoodieCDCLogRecord<T>> {
 
-  private ClosableIterator<IndexedRecord> itr;
-
-  private IndexedRecord record;
-
-  public HoodieCDCLogRecordIterator(HoodieStorage storage, HoodieLogFile[] 
cdcLogFiles, HoodieSchema cdcSchema) {
-    this.storage = storage;
-    this.cdcSchema = cdcSchema;
-    this.cdcLogFileIter = Arrays.stream(cdcLogFiles).iterator();
-  }
-
-  @Override
-  public boolean hasNext() {
-    if (record != null) {
-      return true;
-    }
-    if (itr == null || !itr.hasNext()) {
-      if (reader == null || !reader.hasNext()) {
-        // step1: load new file reader first.
-        if (!loadReader()) {
-          return false;
-        }
-      }
-      // step2: load block records iterator
-      if (!loadItr()) {
+  static <T> HoodieCDCLogRecordIterator<T> empty() {
+    return new HoodieCDCLogRecordIterator<T>() {
+      @Override
+      public boolean hasNext() {
         return false;
       }
-    }
-    record = itr.next();
-    return true;
-  }
 
-  private boolean loadReader() {
-    try {
-      closeReader();
-      if (cdcLogFileIter.hasNext()) {
-        reader = new HoodieLogFileReader(storage, cdcLogFileIter.next(), 
cdcSchema, HoodieLogFileReader.DEFAULT_BUFFER_SIZE);
-        return reader.hasNext();
+      @Override
+      public HoodieCDCLogRecord<T> next() {
+        throw new NoSuchElementException("No CDC log records");
       }
-      return false;
-    } catch (IOException e) {
-      throw new HoodieIOException(e.getMessage());
-    }
-  }
-
-  private boolean loadItr() {
-    HoodieDataBlock dataBlock = (HoodieDataBlock) reader.next();
-    closeItr();
-    // TODO support cdc with spark record.
-    itr = new 
CloseableMappingIterator(dataBlock.getRecordIterator(HoodieRecordType.AVRO), 
record -> ((HoodieAvroIndexedRecord) record).getData());
-    return itr.hasNext();
-  }
 
-  @Override
-  public IndexedRecord next() {
-    IndexedRecord ret = record;
-    record = null;
-    return ret;
-  }
-
-  @Override
-  public void close() {
-    try {
-      closeItr();
-      closeReader();
-    } catch (IOException e) {
-      throw new HoodieIOException(e.getMessage());
-    }
-  }
-
-  // -------------------------------------------------------------------------
-  //  Utilities
-  // -------------------------------------------------------------------------
-
-  private void closeReader() throws IOException {
-    if (reader != null) {
-      reader.close();
-      reader = null;
-    }
-  }
-
-  private void closeItr() {
-    if (itr != null) {
-      itr.close();
-      itr = null;
-    }
+      @Override
+      public void close() {
+        // no-op
+      }
+    };
   }
 }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCNativeLogRecordIterator.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCNativeLogRecordIterator.java
new file mode 100644
index 000000000000..fb4d4bcd995c
--- /dev/null
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/log/HoodieCDCNativeLogRecordIterator.java
@@ -0,0 +1,117 @@
+/*
+ * 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.common.table.log;
+
+import org.apache.hudi.common.util.ValidationUtils;
+import org.apache.hudi.common.util.collection.ClosableIterator;
+
+import org.apache.avro.generic.IndexedRecord;
+
+import java.util.Iterator;
+import java.util.NoSuchElementException;
+import java.util.function.Function;
+
+/**
+ * CDC log record iterator for native CDC log files read as engine-specific 
rows.
+ *
+ * @param <T> Engine-specific record type used by native CDC log files
+ */
+public class HoodieCDCNativeLogRecordIterator<T> implements 
HoodieCDCLogRecordIterator<T> {
+
+  private final Iterator<String> cdcFileIterator;
+  private final Function<String, ClosableIterator<T>> recordIteratorFunc;
+  private final HoodieCDCEngineRecordAccessor<T> recordAccessor;
+  private ClosableIterator<T> recordIterator;
+
+  public HoodieCDCNativeLogRecordIterator(
+      Iterator<String> cdcFileIterator,
+      Function<String, ClosableIterator<T>> recordIteratorFunc,
+      HoodieCDCEngineRecordAccessor<T> recordAccessor) {
+    this.cdcFileIterator = cdcFileIterator;
+    this.recordIteratorFunc = recordIteratorFunc;
+    this.recordAccessor = recordAccessor;
+  }
+
+  @Override
+  public boolean hasNext() {
+    while (recordIterator == null || !recordIterator.hasNext()) {
+      if (recordIterator != null) {
+        recordIterator.close();
+        recordIterator = null;
+      }
+      if (!cdcFileIterator.hasNext()) {
+        return false;
+      }
+      recordIterator = recordIteratorFunc.apply(cdcFileIterator.next());
+      ValidationUtils.checkState(recordIterator != null, "Native CDC record 
iterator must not be null");
+    }
+    return true;
+  }
+
+  @Override
+  public HoodieCDCLogRecord<T> next() {
+    if (!hasNext()) {
+      throw new NoSuchElementException("No more CDC log records");
+    }
+    return new NativeCDCLogRecord<>(recordIterator.next(), recordAccessor);
+  }
+
+  @Override
+  public void close() {
+    if (recordIterator != null) {
+      recordIterator.close();
+      recordIterator = null;
+    }
+  }
+
+  private static class NativeCDCLogRecord<T> implements HoodieCDCLogRecord<T> {
+    private final T record;
+    private final HoodieCDCEngineRecordAccessor<T> recordAccessor;
+
+    private NativeCDCLogRecord(T record, HoodieCDCEngineRecordAccessor<T> 
recordAccessor) {
+      this.record = record;
+      this.recordAccessor = recordAccessor;
+    }
+
+    @Override
+    public String getOperation() {
+      return recordAccessor.getOperation(record);
+    }
+
+    @Override
+    public String getRecordKey() {
+      return recordAccessor.getRecordKey(record);
+    }
+
+    @Override
+    public IndexedRecord getAvroImage(int ordinal) {
+      throw new UnsupportedOperationException("Native CDC records do not 
contain Avro images");
+    }
+
+    @Override
+    public T getEngineImage(int ordinal, int imageArity) {
+      return recordAccessor.getImage(record, ordinal, imageArity);
+    }
+
+    @Override
+    public boolean isNative() {
+      return true;
+    }
+  }
+}
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/log/TestHoodieCDCNativeLogRecordIterator.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/log/TestHoodieCDCNativeLogRecordIterator.java
new file mode 100644
index 000000000000..ec3bff185c6e
--- /dev/null
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/log/TestHoodieCDCNativeLogRecordIterator.java
@@ -0,0 +1,162 @@
+/*
+ * 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.common.table.log;
+
+import org.apache.hudi.common.util.collection.ClosableIterator;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Iterator;
+import java.util.LinkedHashMap;
+import java.util.Map;
+import java.util.NoSuchElementException;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class TestHoodieCDCNativeLogRecordIterator {
+
+  @Test
+  public void testIteratesAcrossFilesAndClosesExhaustedIterators() {
+    Map<String, TrackingClosableIterator<String>> fileIterators = new 
LinkedHashMap<>();
+    fileIterators.put("file1", new 
TrackingClosableIterator<>(Collections.singletonList("i:key1:before1:after1")));
+    fileIterators.put("file2", new 
TrackingClosableIterator<>(Collections.emptyList()));
+    fileIterators.put("file3", new 
TrackingClosableIterator<>(Arrays.asList("u:key2:before2:after2", 
"d:key3:before3:null")));
+
+    HoodieCDCNativeLogRecordIterator<String> iterator = new 
HoodieCDCNativeLogRecordIterator<>(
+        fileIterators.keySet().iterator(),
+        fileIterators::get,
+        accessor());
+
+    assertTrue(iterator.hasNext());
+    HoodieCDCLogRecord<String> insertRecord = iterator.next();
+    assertEquals("i", insertRecord.getOperation());
+    assertEquals("key1", insertRecord.getRecordKey());
+    assertEquals("after1", insertRecord.getEngineImage(3, 2));
+    assertTrue(insertRecord.isNative());
+
+    assertTrue(iterator.hasNext());
+    HoodieCDCLogRecord<String> updateRecord = iterator.next();
+    assertEquals("u", updateRecord.getOperation());
+    assertEquals("key2", updateRecord.getRecordKey());
+    assertEquals("before2", updateRecord.getEngineImage(2, 2));
+
+    assertTrue(iterator.hasNext());
+    HoodieCDCLogRecord<String> deleteRecord = iterator.next();
+    assertEquals("d", deleteRecord.getOperation());
+    assertEquals("key3", deleteRecord.getRecordKey());
+    assertNull(deleteRecord.getEngineImage(3, 2));
+
+    assertFalse(iterator.hasNext());
+    assertTrue(fileIterators.get("file1").isClosed());
+    assertTrue(fileIterators.get("file2").isClosed());
+    assertTrue(fileIterators.get("file3").isClosed());
+  }
+
+  @Test
+  public void testCloseClosesCurrentIterator() {
+    Map<String, TrackingClosableIterator<String>> fileIterators = new 
LinkedHashMap<>();
+    fileIterators.put("file1", new 
TrackingClosableIterator<>(Arrays.asList("i:key1:before1:after1", 
"u:key2:before2:after2")));
+
+    HoodieCDCNativeLogRecordIterator<String> iterator = new 
HoodieCDCNativeLogRecordIterator<>(
+        fileIterators.keySet().iterator(),
+        fileIterators::get,
+        accessor());
+
+    assertTrue(iterator.hasNext());
+    iterator.close();
+
+    assertTrue(fileIterators.get("file1").isClosed());
+  }
+
+  @Test
+  public void testRejectsNullNativeRecordIterator() {
+    HoodieCDCNativeLogRecordIterator<String> iterator = new 
HoodieCDCNativeLogRecordIterator<>(
+        Collections.singletonList("file1").iterator(),
+        cdcFile -> null,
+        accessor());
+
+    assertThrows(IllegalStateException.class, iterator::hasNext);
+  }
+
+  @Test
+  public void testEmptyIteratorThrowsOnNext() {
+    HoodieCDCLogRecordIterator<String> iterator = 
HoodieCDCLogRecordIterator.empty();
+
+    assertFalse(iterator.hasNext());
+    assertThrows(NoSuchElementException.class, iterator::next);
+  }
+
+  private static HoodieCDCEngineRecordAccessor<String> accessor() {
+    return new HoodieCDCEngineRecordAccessor<String>() {
+      @Override
+      public String getOperation(String record) {
+        return split(record)[0];
+      }
+
+      @Override
+      public String getRecordKey(String record) {
+        return split(record)[1];
+      }
+
+      @Override
+      public String getImage(String record, int ordinal, int imageArity) {
+        String image = split(record)[ordinal];
+        return "null".equals(image) ? null : image;
+      }
+    };
+  }
+
+  private static String[] split(String record) {
+    return record.split(":");
+  }
+
+  private static class TrackingClosableIterator<T> implements 
ClosableIterator<T> {
+    private final Iterator<T> iterator;
+    private boolean closed;
+
+    private TrackingClosableIterator(Iterable<T> records) {
+      this.iterator = records.iterator();
+    }
+
+    @Override
+    public boolean hasNext() {
+      return iterator.hasNext();
+    }
+
+    @Override
+    public T next() {
+      return iterator.next();
+    }
+
+    @Override
+    public void close() {
+      closed = true;
+    }
+
+    private boolean isClosed() {
+      return closed;
+    }
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/function/HoodieCdcSplitReaderFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/function/HoodieCdcSplitReaderFunction.java
index fd9e2f6fca54..b539bc835f14 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/function/HoodieCdcSplitReaderFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/function/HoodieCdcSplitReaderFunction.java
@@ -210,16 +210,16 @@ public class HoodieCdcSplitReaderFunction extends 
AbstractSplitReaderFunction {
         switch (mode) {
           case DATA_BEFORE_AFTER:
             return new CdcIterators.BeforeAfterImageIterator(
-                getHadoopConf(), tablePath, tableSchema, requiredSchema,
+                conf, getHadoopConf(), tablePath, tableSchema, requiredSchema,
                 tableState.getRequiredRowType(), cdcSchema, fileSplit);
           case DATA_BEFORE:
             return new CdcIterators.BeforeImageIterator(
-                getHadoopConf(), tablePath, tableSchema, requiredSchema,
+                conf, getHadoopConf(), tablePath, tableSchema, requiredSchema,
                 tableState.getRequiredRowType(), 
tableState.getRequiredPositions(),
                 maxCompactionMemoryInBytes, cdcSchema, fileSplit, 
imageManager);
           case OP_KEY_ONLY:
             return new CdcIterators.RecordKeyImageIterator(
-                getHadoopConf(), tablePath, tableSchema, requiredSchema,
+                conf, getHadoopConf(), tablePath, tableSchema, requiredSchema,
                 tableState.getRequiredRowType(), 
tableState.getRequiredPositions(),
                 maxCompactionMemoryInBytes, cdcSchema, fileSplit, 
imageManager);
           default:
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcInputFormat.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcInputFormat.java
index 27bb2b3a56e8..af28d2e249a0 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcInputFormat.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcInputFormat.java
@@ -138,14 +138,14 @@ public class CdcInputFormat extends 
MergeOnReadInputFormat {
         switch (mode) {
           case DATA_BEFORE_AFTER:
             return new CdcIterators.BeforeAfterImageIterator(
-                hadoopConf, tablePath, tblSchema, reqSchema, 
tableState.getRequiredRowType(), cdcSchema, fileSplit);
+                conf, hadoopConf, tablePath, tblSchema, reqSchema, 
tableState.getRequiredRowType(), cdcSchema, fileSplit);
           case DATA_BEFORE:
             return new CdcIterators.BeforeImageIterator(
-                hadoopConf, tablePath, tblSchema, reqSchema, 
tableState.getRequiredRowType(),
+                conf, hadoopConf, tablePath, tblSchema, reqSchema, 
tableState.getRequiredRowType(),
                 tableState.getRequiredPositions(), maxCompactionMemoryInBytes, 
cdcSchema, fileSplit, imageManager);
           case OP_KEY_ONLY:
             return new CdcIterators.RecordKeyImageIterator(
-                hadoopConf, tablePath, tblSchema, reqSchema, 
tableState.getRequiredRowType(),
+                conf, hadoopConf, tablePath, tblSchema, reqSchema, 
tableState.getRequiredRowType(),
                 tableState.getRequiredPositions(), maxCompactionMemoryInBytes, 
cdcSchema, fileSplit, imageManager);
           default:
             throw new AssertionError("Unexpected mode: " + mode);
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
index fc688936981a..c4537eaadb5a 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.model.BaseFile;
 import org.apache.hudi.common.model.FileSlice;
+import org.apache.hudi.common.model.HoodieFileFormat;
 import org.apache.hudi.common.model.HoodieLogFile;
 import org.apache.hudi.common.model.HoodieOperation;
 import org.apache.hudi.common.model.HoodieRecord;
@@ -33,7 +34,11 @@ import org.apache.hudi.common.schema.HoodieSchemaUtils;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit;
 import org.apache.hudi.common.table.cdc.HoodieCDCUtils;
+import org.apache.hudi.common.table.log.HoodieCDCEngineRecordAccessor;
+import org.apache.hudi.common.table.log.HoodieCDCInlineLogRecordIterator;
+import org.apache.hudi.common.table.log.HoodieCDCLogRecord;
 import org.apache.hudi.common.table.log.HoodieCDCLogRecordIterator;
+import org.apache.hudi.common.table.log.HoodieCDCNativeLogRecordIterator;
 import org.apache.hudi.common.table.read.BufferedRecord;
 import org.apache.hudi.common.table.read.BufferedRecordMerger;
 import org.apache.hudi.common.table.read.BufferedRecordMergerFactory;
@@ -50,13 +55,17 @@ import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.configuration.FlinkOptions;
 import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.hadoop.fs.HadoopFSUtils;
+import org.apache.hudi.io.storage.HoodieIOFactory;
 import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.HoodieStorageUtils;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.format.FlinkReaderContextFactory;
 import org.apache.hudi.table.format.FormatUtils;
+import org.apache.hudi.table.format.HoodieRowDataFileReader;
+import org.apache.hudi.table.format.InternalSchemaManager;
 import org.apache.hudi.table.format.mor.MergeOnReadInputSplit;
 import org.apache.hudi.util.AvroToRowDataConverters;
+import org.apache.hudi.util.FlinkWriteClients;
 import org.apache.hudi.util.HoodieSchemaConverter;
 import org.apache.hudi.util.RowDataProjection;
 
@@ -364,21 +373,42 @@ public final class CdcIterators {
 
   /**
    * Base iterator for CDC log files stored with supplemental logging (AS_IS 
inference case).
-   * Reads a {@link HoodieCDCLogRecordIterator} and resolves before/after 
images using
-   * subclass-specific logic.
+   * Reads inline or native CDC log files through a normalized CDC record 
iterator and resolves
+   * before/after images using subclass-specific logic.
    */
   public abstract static class BaseImageIterator implements 
ClosableIterator<RowData> {
     private final HoodieSchema requiredSchema;
     private final int[] requiredPos;
     private final GenericRecordBuilder recordBuilder;
     private final AvroToRowDataConverters.AvroToRowDataConverter 
avroToRowDataConverter;
-    private HoodieCDCLogRecordIterator cdcItr;
+    private final RowDataProjection nativeCdcImageProjection;
+    private final int nativeCdcImageArity;
+    private static final HoodieCDCEngineRecordAccessor<RowData> 
ROW_DATA_CDC_RECORD_ACCESSOR =
+        new HoodieCDCEngineRecordAccessor<RowData>() {
+          @Override
+          public String getOperation(RowData record) {
+            return record.getString(0).toString();
+          }
+
+          @Override
+          public String getRecordKey(RowData record) {
+            return record.getString(1).toString();
+          }
+
+          @Override
+          public RowData getImage(RowData record, int ordinal, int imageArity) 
{
+            return record.isNullAt(ordinal) ? null : record.getRow(ordinal, 
imageArity);
+          }
+        };
 
-    private GenericRecord cdcRecord;
+    private HoodieCDCLogRecordIterator<?> cdcItr;
+
+    private HoodieCDCLogRecord<?> cdcRecord;
     private RowData sideImage;
     private RowData currentImage;
 
     protected BaseImageIterator(
+        org.apache.flink.configuration.Configuration conf,
         org.apache.hadoop.conf.Configuration hadoopConf,
         String tablePath,
         HoodieSchema tableSchema,
@@ -390,20 +420,98 @@ public final class CdcIterators {
       this.requiredPos = computeRequiredPos(tableSchema, requiredSchema);
       this.recordBuilder = new 
GenericRecordBuilder(requiredSchema.getAvroSchema());
       this.avroToRowDataConverter = 
AvroToRowDataConverters.createRowConverter(requiredSchema, requiredRowType, 
true);
+      this.nativeCdcImageProjection = 
RowDataProjection.instance(requiredRowType, requiredPos);
+      this.nativeCdcImageArity = 
HoodieSchemaUtils.removeMetadataFields(tableSchema).getFields().size();
+      this.cdcItr = createCdcRecordIterator(
+          conf, hadoopConf, tablePath, cdcSchema, fileSplit);
+    }
 
-      StoragePath hadoopTablePath = new StoragePath(tablePath);
-      HoodieStorage storage = HoodieStorageUtils.getStorage(
-          tablePath, HadoopFSUtils.getStorageConf(hadoopConf));
-      HoodieLogFile[] cdcLogFiles = fileSplit.getCdcFiles().stream()
-          .map(cdcFile -> {
-            try {
-              return new HoodieLogFile(storage.getPathInfo(new 
StoragePath(hadoopTablePath, cdcFile)));
-            } catch (IOException e) {
-              throw new HoodieIOException("Failed to get file status for CDC 
log: " + cdcFile, e);
-            }
-          })
-          .toArray(HoodieLogFile[]::new);
-      this.cdcItr = new HoodieCDCLogRecordIterator(storage, cdcLogFiles, 
cdcSchema);
+    private static HoodieCDCLogRecordIterator<?> createCdcRecordIterator(
+        org.apache.flink.configuration.Configuration conf,
+        org.apache.hadoop.conf.Configuration hadoopConf,
+        String tablePath,
+        HoodieSchema cdcSchema,
+        HoodieCDCFileSplit fileSplit) {
+      if (fileSplit.getCdcFiles() == null || 
fileSplit.getCdcFiles().isEmpty()) {
+        return HoodieCDCLogRecordIterator.empty();
+      }
+      if (isNativeCdcFileSplit(fileSplit)) {
+        return new HoodieCDCNativeLogRecordIterator<>(
+            fileSplit.getCdcFiles().iterator(),
+            cdcFile -> getNativeCdcFileIterator(conf, hadoopConf, tablePath, 
cdcFile, cdcSchema),
+            ROW_DATA_CDC_RECORD_ACCESSOR);
+      } else {
+        StoragePath hadoopTablePath = new StoragePath(tablePath);
+        HoodieStorage storage = HoodieStorageUtils.getStorage(
+            tablePath, HadoopFSUtils.getStorageConf(hadoopConf));
+        HoodieLogFile[] cdcLogFiles = fileSplit.getCdcFiles().stream()
+            .map(cdcFile -> {
+              try {
+                return new HoodieLogFile(storage.getPathInfo(new 
StoragePath(hadoopTablePath, cdcFile)));
+              } catch (IOException e) {
+                throw new HoodieIOException("Failed to get file status for CDC 
log: " + cdcFile, e);
+              }
+            })
+            .toArray(HoodieLogFile[]::new);
+        return new HoodieCDCInlineLogRecordIterator(storage, cdcLogFiles, 
cdcSchema);
+      }
+    }
+
+    private static ClosableIterator<RowData> getNativeCdcFileIterator(
+        org.apache.flink.configuration.Configuration conf,
+        org.apache.hadoop.conf.Configuration hadoopConf,
+        String tablePath,
+        String cdcFile,
+        HoodieSchema cdcSchema) {
+      HoodieRowDataFileReader reader = null;
+      try {
+        StoragePath cdcFilePath = new StoragePath(tablePath, cdcFile);
+        HoodieStorage storage = HoodieStorageUtils.getStorage(
+            tablePath, HadoopFSUtils.getStorageConf(hadoopConf));
+        HoodieFileFormat cdcFileFormat = 
HoodieFileFormat.fromFileExtension(cdcFilePath.getFileExtension());
+        reader = (HoodieRowDataFileReader) 
HoodieIOFactory.getIOFactory(storage)
+            .getReaderFactory(HoodieRecord.HoodieRecordType.FLINK)
+            .getFileReader(
+                FlinkWriteClients.getHoodieClientConfig(conf), cdcFilePath, 
cdcFileFormat, Option.empty());
+        return closeReaderWithIterator(
+            reader,
+            reader.getRowDataIterator(cdcSchema, cdcSchema, 
InternalSchemaManager.DISABLED, Collections.emptyList()));
+      } catch (IOException e) {
+        if (reader != null) {
+          reader.close();
+        }
+        throw new HoodieIOException("Failed to create native CDC record 
iterator for file: " + cdcFile, e);
+      } catch (RuntimeException e) {
+        if (reader != null) {
+          reader.close();
+        }
+        throw e;
+      }
+    }
+
+    private static ClosableIterator<RowData> closeReaderWithIterator(
+        HoodieRowDataFileReader reader,
+        ClosableIterator<RowData> iterator) {
+      return new ClosableIterator<RowData>() {
+        @Override
+        public boolean hasNext() {
+          return iterator.hasNext();
+        }
+
+        @Override
+        public RowData next() {
+          return iterator.next();
+        }
+
+        @Override
+        public void close() {
+          try {
+            iterator.close();
+          } finally {
+            reader.close();
+          }
+        }
+      };
     }
 
     private static int[] computeRequiredPos(HoodieSchema tableSchema, 
HoodieSchema requiredSchema) {
@@ -417,6 +525,14 @@ public final class CdcIterators {
           .toArray();
     }
 
+    private static boolean isNativeCdcFileSplit(HoodieCDCFileSplit fileSplit) {
+      boolean nativeCdc = 
FSUtils.matchNativeLogFile(fileSplit.getCdcFiles().get(0)).isPresent();
+      ValidationUtils.checkState(fileSplit.getCdcFiles().stream()
+              .allMatch(path -> FSUtils.matchNativeLogFile(path).isPresent() 
== nativeCdc),
+          "CDC file split cannot mix inline and native CDC log files");
+      return nativeCdc;
+    }
+
     @Override
     public boolean hasNext() {
       if (sideImage != null) {
@@ -424,17 +540,16 @@ public final class CdcIterators {
         sideImage = null;
         return true;
       } else if (cdcItr.hasNext()) {
-        cdcRecord = (GenericRecord) cdcItr.next();
-        String op = String.valueOf(cdcRecord.get(0));
-        resolveImage(op);
+        cdcRecord = cdcItr.next();
+        resolveImage(cdcRecord.getOperation());
         return true;
       }
       return false;
     }
 
-    protected abstract RowData getAfterImage(RowKind rowKind, GenericRecord 
cdcRecord);
+    protected abstract RowData getAfterImage(RowKind rowKind, 
HoodieCDCLogRecord<?> cdcRecord);
 
-    protected abstract RowData getBeforeImage(RowKind rowKind, GenericRecord 
cdcRecord);
+    protected abstract RowData getBeforeImage(RowKind rowKind, 
HoodieCDCLogRecord<?> cdcRecord);
 
     @Override
     public RowData next() {
@@ -473,6 +588,18 @@ public final class CdcIterators {
       resolved.setRowKind(rowKind);
       return resolved;
     }
+
+    protected RowData resolveImage(RowKind rowKind, HoodieCDCLogRecord<?> 
cdcRecord, int ordinal) {
+      if (!cdcRecord.isNative()) {
+        return resolveAvro(rowKind, (GenericRecord) 
cdcRecord.getAvroImage(ordinal));
+      }
+      RowData image = (RowData) cdcRecord.getEngineImage(ordinal, 
nativeCdcImageArity);
+      if (image == null) {
+        return null;
+      }
+      image.setRowKind(rowKind);
+      return nativeCdcImageProjection.project(image);
+    }
   }
 
   /**
@@ -481,6 +608,7 @@ public final class CdcIterators {
    */
   public static class BeforeAfterImageIterator extends BaseImageIterator {
     public BeforeAfterImageIterator(
+        org.apache.flink.configuration.Configuration conf,
         org.apache.hadoop.conf.Configuration hadoopConf,
         String tablePath,
         HoodieSchema tableSchema,
@@ -488,17 +616,17 @@ public final class CdcIterators {
         RowType requiredRowType,
         HoodieSchema cdcSchema,
         HoodieCDCFileSplit fileSplit) {
-      super(hadoopConf, tablePath, tableSchema, requiredSchema, 
requiredRowType, cdcSchema, fileSplit);
+      super(conf, hadoopConf, tablePath, tableSchema, requiredSchema, 
requiredRowType, cdcSchema, fileSplit);
     }
 
     @Override
-    protected RowData getAfterImage(RowKind rowKind, GenericRecord cdcRecord) {
-      return resolveAvro(rowKind, (GenericRecord) cdcRecord.get(3));
+    protected RowData getAfterImage(RowKind rowKind, HoodieCDCLogRecord<?> 
cdcRecord) {
+      return resolveImage(rowKind, cdcRecord, 3);
     }
 
     @Override
-    protected RowData getBeforeImage(RowKind rowKind, GenericRecord cdcRecord) 
{
-      return resolveAvro(rowKind, (GenericRecord) cdcRecord.get(2));
+    protected RowData getBeforeImage(RowKind rowKind, HoodieCDCLogRecord<?> 
cdcRecord) {
+      return resolveImage(rowKind, cdcRecord, 2);
     }
   }
 
@@ -514,6 +642,7 @@ public final class CdcIterators {
     protected final CdcImageManager imageManager;
 
     public BeforeImageIterator(
+        org.apache.flink.configuration.Configuration conf,
         org.apache.hadoop.conf.Configuration hadoopConf,
         String tablePath,
         HoodieSchema tableSchema,
@@ -524,7 +653,7 @@ public final class CdcIterators {
         HoodieSchema cdcSchema,
         HoodieCDCFileSplit fileSplit,
         CdcImageManager imageManager) throws IOException {
-      super(hadoopConf, tablePath, tableSchema, requiredSchema, 
requiredRowType, cdcSchema, fileSplit);
+      super(conf, hadoopConf, tablePath, tableSchema, requiredSchema, 
requiredRowType, cdcSchema, fileSplit);
       this.maxCompactionMemoryInBytes = maxCompactionMemoryInBytes;
       this.projection = RowDataProjection.instance(requiredRowType, 
requiredPositions);
       this.imageManager = imageManager;
@@ -539,16 +668,16 @@ public final class CdcIterators {
     }
 
     @Override
-    protected RowData getAfterImage(RowKind rowKind, GenericRecord cdcRecord) {
-      String recordKey = cdcRecord.get(1).toString();
+    protected RowData getAfterImage(RowKind rowKind, HoodieCDCLogRecord<?> 
cdcRecord) {
+      String recordKey = cdcRecord.getRecordKey();
       RowData row = imageManager.getImageRecord(recordKey, afterImages, 
rowKind);
       row.setRowKind(rowKind);
       return projection.project(row);
     }
 
     @Override
-    protected RowData getBeforeImage(RowKind rowKind, GenericRecord cdcRecord) 
{
-      return resolveAvro(rowKind, (GenericRecord) cdcRecord.get(2));
+    protected RowData getBeforeImage(RowKind rowKind, HoodieCDCLogRecord<?> 
cdcRecord) {
+      return resolveImage(rowKind, cdcRecord, 2);
     }
   }
 
@@ -561,6 +690,7 @@ public final class CdcIterators {
     protected ExternalSpillableMap<String, byte[]> beforeImages;
 
     public RecordKeyImageIterator(
+        org.apache.flink.configuration.Configuration conf,
         org.apache.hadoop.conf.Configuration hadoopConf,
         String tablePath,
         HoodieSchema tableSchema,
@@ -571,7 +701,7 @@ public final class CdcIterators {
         HoodieSchema cdcSchema,
         HoodieCDCFileSplit fileSplit,
         CdcImageManager imageManager) throws IOException {
-      super(hadoopConf, tablePath, tableSchema, requiredSchema, 
requiredRowType,
+      super(conf, hadoopConf, tablePath, tableSchema, requiredSchema, 
requiredRowType,
           requiredPositions, maxCompactionMemoryInBytes, cdcSchema, fileSplit, 
imageManager);
     }
 
@@ -585,8 +715,8 @@ public final class CdcIterators {
     }
 
     @Override
-    protected RowData getBeforeImage(RowKind rowKind, GenericRecord cdcRecord) 
{
-      String recordKey = cdcRecord.get(1).toString();
+    protected RowData getBeforeImage(RowKind rowKind, HoodieCDCLogRecord<?> 
cdcRecord) {
+      String recordKey = cdcRecord.getRecordKey();
       RowData row = imageManager.getImageRecord(recordKey, beforeImages, 
rowKind);
       row.setRowKind(rowKind);
       return projection.project(row);
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
index cc316243d266..04df26a839ef 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala
@@ -36,11 +36,11 @@ import 
org.apache.hudi.common.table.cdc.{HoodieCDCFileSplit, HoodieCDCUtils}
 import org.apache.hudi.common.table.cdc.HoodieCDCInferenceCase._
 import org.apache.hudi.common.table.cdc.HoodieCDCOperation._
 import org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode._
-import org.apache.hudi.common.table.log.{HoodieCDCLogRecordIterator, 
HoodieMergedLogRecordReader}
+import org.apache.hudi.common.table.log.{HoodieCDCEngineRecordAccessor, 
HoodieCDCInlineLogRecordIterator, HoodieCDCLogRecord, 
HoodieCDCLogRecordIterator, HoodieCDCNativeLogRecordIterator, 
HoodieMergedLogRecordReader}
 import org.apache.hudi.common.table.read.{BufferedRecord, 
BufferedRecordMerger, BufferedRecordMergerFactory, BufferedRecords, 
FileGroupReaderSchemaHandler, HoodieFileGroupReader, HoodieReadStats, 
IteratorMode, UpdateProcessor}
 import org.apache.hudi.common.table.read.buffer.KeyBasedFileGroupRecordBuffer
-import org.apache.hudi.common.util.{DefaultSizeEstimator, HoodieRecordUtils, 
Option}
-import org.apache.hudi.common.util.collection.ExternalSpillableMap
+import org.apache.hudi.common.util.{DefaultSizeEstimator, HoodieRecordUtils, 
Option, ValidationUtils}
+import org.apache.hudi.common.util.collection.{ClosableIterator, 
ExternalSpillableMap}
 import org.apache.hudi.config.HoodieWriteConfig
 import org.apache.hudi.data.CloseableIteratorListener
 import org.apache.hudi.io.util.FileIOUtils
@@ -52,7 +52,6 @@ import 
org.apache.parquet.avro.HoodieAvroParquetSchemaConverter.getAvroSchemaCon
 import org.apache.spark.Partition
 import 
org.apache.spark.sql.HoodieCatalystExpressionUtils.generateUnsafeProjection
 import org.apache.spark.sql.HoodieInternalRowUtils
-import org.apache.spark.sql.avro.HoodieAvroDeserializer
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.catalyst.expressions.Projection
 import org.apache.spark.sql.execution.datasources.SparkColumnarFileReader
@@ -153,13 +152,6 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
     
org.apache.hudi.common.util.Option.empty[org.apache.parquet.schema.MessageType]()
   }
 
-  /**
-   * The deserializer used to convert the CDC GenericRecord to Spark 
InternalRow.
-   */
-  private lazy val cdcRecordDeserializer: HoodieAvroDeserializer = {
-    sparkAdapter.createAvroDeserializer(cdcHoodieSchema, cdcSparkSchema)
-  }
-
   private lazy val projection: Projection = 
generateUnsafeProjection(cdcSchema, requiredCdcSchema)
 
   // Iterator on cdc file
@@ -189,7 +181,7 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
   /**
    * Only one case where it will be used is that extract the change data from 
cdc log files.
    */
-  private var cdcLogRecordIterator: HoodieCDCLogRecordIterator = _
+  private var cdcRecordIterator: HoodieCDCLogRecordIterator[_] = 
HoodieCDCLogRecordIterator.empty()
 
   /**
    * The next record need to be returned when call next().
@@ -233,10 +225,21 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
   // images. Keyed by the record's schema id so schema evolution is handled 
correctly.
   private val cdcImageConverterMap: mutable.Map[Integer, 
(UnaryOperator[InternalRow], InternalRowToJsonStringConverter)] = 
mutable.Map.empty
 
+  private lazy val cdcDataSparkSchema: StructType = 
HoodieSchemaConversionUtils.convertHoodieSchemaToStructType(
+    HoodieSchemaUtils.removeMetadataFields(schema))
+
+  private lazy val nativeCdcParquetSchemaOpt = {
+    val hadoopConf = storage.getConf.unwrapAs(classOf[Configuration])
+    val parquetSchema = 
getAvroSchemaConverter(hadoopConf).convert(cdcHoodieSchema)
+    org.apache.hudi.common.util.Option.of(parquetSchema)
+  }
+
+  private lazy val nativeCdcImageConverter = new 
InternalRowToJsonStringConverter(cdcDataSparkSchema)
+
   private def needLoadNextFile: Boolean = {
     !recordIter.hasNext &&
       !logRecordIter.hasNext &&
-      (cdcLogRecordIterator == null || !cdcLogRecordIterator.hasNext)
+      !cdcRecordIterator.hasNext
   }
 
   @tailrec final def hasNextInternal: Boolean = {
@@ -260,7 +263,7 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
             hasNextInternal
           }
         case AS_IS =>
-          if (cdcLogRecordIterator.hasNext && loadNext()) {
+          if (cdcRecordIterator.hasNext && loadNext()) {
             true
           } else {
             hasNextInternal
@@ -296,46 +299,7 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
       case LOG_FILE =>
         loaded = loadNextLogRecord()
       case AS_IS =>
-        val record = cdcLogRecordIterator.next().asInstanceOf[GenericRecord]
-        cdcSupplementalLoggingMode match {
-          case `DATA_BEFORE_AFTER` =>
-            recordToLoad.update(0, 
convertToUTF8String(String.valueOf(record.get(0))))
-            val before = record.get(2).asInstanceOf[GenericRecord]
-            recordToLoad.update(2, recordToJsonAsUTF8String(before))
-            val after = record.get(3).asInstanceOf[GenericRecord]
-            recordToLoad.update(3, recordToJsonAsUTF8String(after))
-          case `DATA_BEFORE` =>
-            val row = 
cdcRecordDeserializer.deserialize(record).get.asInstanceOf[InternalRow]
-            val op = row.getString(0)
-            val recordKey = row.getString(1)
-            recordToLoad.update(0, convertToUTF8String(op))
-            val before = record.get(2).asInstanceOf[GenericRecord]
-            recordToLoad.update(2, recordToJsonAsUTF8String(before))
-            parse(op) match {
-              case INSERT =>
-                recordToLoad.update(3, 
convertBufferedRecordToJsonString(afterImageRecords.get(recordKey)))
-              case UPDATE =>
-                recordToLoad.update(3, 
convertBufferedRecordToJsonString(afterImageRecords.get(recordKey)))
-              case _ =>
-                recordToLoad.update(3, null)
-            }
-          case _ =>
-            val row = 
cdcRecordDeserializer.deserialize(record).get.asInstanceOf[InternalRow]
-            val op = row.getString(0)
-            val recordKey = row.getString(1)
-            recordToLoad.update(0, convertToUTF8String(op))
-            parse(op) match {
-              case INSERT =>
-                recordToLoad.update(2, null)
-                recordToLoad.update(3, 
convertBufferedRecordToJsonString(afterImageRecords.get(recordKey)))
-              case UPDATE =>
-                recordToLoad.update(2, 
convertBufferedRecordToJsonString(beforeImageRecords(recordKey)))
-                recordToLoad.update(3, 
convertBufferedRecordToJsonString(afterImageRecords.get(recordKey)))
-              case _ =>
-                recordToLoad.update(2, 
convertBufferedRecordToJsonString(beforeImageRecords(recordKey)))
-                recordToLoad.update(3, null)
-            }
-        }
+        loadNextCdcRecord()
         loaded = true
       case REPLACE_COMMIT =>
         val originRecord = recordIter.next()
@@ -395,14 +359,12 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
     // reset all the iterator first.
     recordIter = Iterator.empty
     logRecordIter = Iterator.empty
+    cdcRecordIterator.close()
+    cdcRecordIterator = HoodieCDCLogRecordIterator.empty()
     keyBasedFileGroupRecordBuffer.ifPresent(k => k.close())
     keyBasedFileGroupRecordBuffer = 
Option.empty.asInstanceOf[Option[KeyBasedFileGroupRecordBuffer[InternalRow]]]
     beforeImageRecords.clear()
     afterImageRecords.clear()
-    if (cdcLogRecordIterator != null) {
-      cdcLogRecordIterator.close()
-      cdcLogRecordIterator = null
-    }
 
     if (cdcFileIter.hasNext) {
       val split = cdcFileIter.next()
@@ -444,10 +406,7 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
             }
           }
 
-          val cdcLogFiles = currentCDCFileSplit.getCdcFiles.asScala.map { 
cdcFile =>
-            new HoodieLogFile(storage.getPathInfo(new StoragePath(basePath, 
cdcFile)))
-          }.toArray
-          cdcLogRecordIterator = new HoodieCDCLogRecordIterator(storage, 
cdcLogFiles, cdcHoodieSchema)
+          cdcRecordIterator = createCdcRecordIterator(currentCDCFileSplit)
         case REPLACE_COMMIT =>
           if (currentCDCFileSplit.getBeforeFileSlice.isPresent) {
             
loadBeforeFileSliceIfNeeded(currentCDCFileSplit.getBeforeFileSlice.get)
@@ -558,6 +517,91 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
     
CloseableIteratorListener.addListener(keyBasedFileGroupRecordBuffer.get().getLogRecordIterator).asScala
   }
 
+  private def readNativeCdcFile(cdcFile: String): Iterator[InternalRow] = {
+    val absCDCPath = new StoragePath(basePath, cdcFile)
+    val fileStatus = storage.getPathInfo(absCDCPath)
+    val pf = sparkPartitionedFileUtils.createPartitionedFile(
+      InternalRow.empty, absCDCPath, 0, fileStatus.getLength)
+    baseFileReader.read(pf, cdcSparkSchema, new StructType(),
+      org.apache.hudi.common.util.Option.empty(), Seq.empty, conf, 
nativeCdcParquetSchemaOpt)
+  }
+
+  private val nativeCdcRecordAccessor = new 
HoodieCDCEngineRecordAccessor[InternalRow] {
+    override def getOperation(record: InternalRow): String = 
record.getString(0)
+    override def getRecordKey(record: InternalRow): String = 
record.getString(1)
+    override def getImage(record: InternalRow, ordinal: Int, imageArity: Int): 
InternalRow = {
+      if (record.isNullAt(ordinal)) null else record.getStruct(ordinal, 
imageArity)
+    }
+  }
+
+  private def createCdcRecordIterator(fileSplit: HoodieCDCFileSplit): 
HoodieCDCLogRecordIterator[_] = {
+    if (fileSplit.getCdcFiles == null || fileSplit.getCdcFiles.isEmpty) {
+      HoodieCDCLogRecordIterator.empty()
+    } else if (isNativeCdcFileSplit(fileSplit)) {
+      new HoodieCDCNativeLogRecordIterator[InternalRow](
+        fileSplit.getCdcFiles.iterator(),
+        cdcFile => ClosableIterator.wrap(readNativeCdcFile(cdcFile).asJava),
+        nativeCdcRecordAccessor)
+    } else {
+      val cdcLogFiles = fileSplit.getCdcFiles.asScala.map { cdcFile =>
+        new HoodieLogFile(storage.getPathInfo(new StoragePath(basePath, 
cdcFile)))
+      }.toArray
+      new HoodieCDCInlineLogRecordIterator(storage, cdcLogFiles, 
cdcHoodieSchema)
+    }
+  }
+
+  private def isNativeCdcFileSplit(fileSplit: HoodieCDCFileSplit): Boolean = {
+    val nativeFlags = fileSplit.getCdcFiles.asScala.map(path => 
FSUtils.matchNativeLogFile(path).isPresent)
+    ValidationUtils.checkState(nativeFlags.forall(_ == nativeFlags.head),
+      "CDC file split cannot mix inline and native CDC log files")
+    nativeFlags.head
+  }
+
+  private def loadNextCdcRecord(): Unit = {
+    val record = cdcRecordIterator.next()
+    cdcSupplementalLoggingMode match {
+      case `DATA_BEFORE_AFTER` =>
+        recordToLoad.update(0, convertToUTF8String(record.getOperation))
+        recordToLoad.update(2, cdcRecordImageToJson(record, 2))
+        recordToLoad.update(3, cdcRecordImageToJson(record, 3))
+      case `DATA_BEFORE` =>
+        recordToLoad.update(0, convertToUTF8String(record.getOperation))
+        recordToLoad.update(2, cdcRecordImageToJson(record, 2))
+        parse(record.getOperation) match {
+          case INSERT | UPDATE =>
+            recordToLoad.update(3, 
convertBufferedRecordToJsonString(afterImageRecords.get(record.getRecordKey)))
+          case _ =>
+            recordToLoad.update(3, null)
+        }
+      case _ =>
+        loadNextKeyOnlyCdcRecord(record)
+    }
+  }
+
+  private def loadNextKeyOnlyCdcRecord(record: HoodieCDCLogRecord[_]): Unit = {
+    recordToLoad.update(0, convertToUTF8String(record.getOperation))
+    parse(record.getOperation) match {
+      case INSERT =>
+        recordToLoad.update(2, null)
+        recordToLoad.update(3, 
convertBufferedRecordToJsonString(afterImageRecords.get(record.getRecordKey)))
+      case UPDATE =>
+        recordToLoad.update(2, 
convertBufferedRecordToJsonString(beforeImageRecords(record.getRecordKey)))
+        recordToLoad.update(3, 
convertBufferedRecordToJsonString(afterImageRecords.get(record.getRecordKey)))
+      case _ =>
+        recordToLoad.update(2, 
convertBufferedRecordToJsonString(beforeImageRecords(record.getRecordKey)))
+        recordToLoad.update(3, null)
+    }
+  }
+
+  private def cdcRecordImageToJson(record: HoodieCDCLogRecord[_], ordinal: 
Int): UTF8String = {
+    if (record.isNative) {
+      val image = record.getEngineImage(ordinal, 
cdcDataSparkSchema.length).asInstanceOf[InternalRow]
+      if (image == null) null else nativeCdcImageConverter.convert(image)
+    } else {
+      
recordToJsonAsUTF8String(record.getAvroImage(ordinal).asInstanceOf[GenericRecord])
+    }
+  }
+
   /**
    * Convert InternalRow to json string.
    */
@@ -586,7 +630,11 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
   }
 
   private def recordToJsonAsUTF8String(record: GenericRecord): UTF8String = {
-    convertToUTF8String(HoodieCDCUtils.recordToJson(record))
+    if (record == null) {
+      null
+    } else {
+      convertToUTF8String(HoodieCDCUtils.recordToJson(record))
+    }
   }
 
   private def merge(currentRecord: BufferedRecord[InternalRow], newRecord: 
BufferedRecord[InternalRow]): BufferedRecord[InternalRow] = {
@@ -605,15 +653,14 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit,
   override def close(): Unit = {
     recordIter = Iterator.empty
     logRecordIter = Iterator.empty
+    cdcRecordIterator.close()
+    cdcRecordIterator = HoodieCDCLogRecordIterator.empty()
     keyBasedFileGroupRecordBuffer.ifPresent(k => k.close())
     keyBasedFileGroupRecordBuffer = 
Option.empty.asInstanceOf[Option[KeyBasedFileGroupRecordBuffer[InternalRow]]]
     beforeImageRecords.clear()
     afterImageRecords.clear()
-    if (cdcLogRecordIterator != null) {
-      cdcLogRecordIterator.close()
-      cdcLogRecordIterator = null
-    }
   }
+
 }
 
 object CDCFileGroupIterator {

Reply via email to