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