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