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