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 {