cshuo commented on code in PR #19307:
URL: https://github.com/apache/hudi/pull/19307#discussion_r3619053364
##########
hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java:
##########
@@ -314,32 +315,85 @@ private ClosableIterator<BufferedRecord<T>>
createFileIterator(StoragePathInfo p
if (FSUtils.isNativeDeleteLogFile(storagePath.getName())) {
return createNativeDeleteLogIterator(pathInfo, storagePath, fileSize);
}
- Pair<HoodieSchema, Map<String, String>> requiredSchemaAndRenamedColumns =
-
readerContext.getSchemaHandler().getRequiredSchemaForFileAndRenamedColumns(storagePath);
- HoodieSchema fileRequiredSchema =
requiredSchemaAndRenamedColumns.getLeft();
- ClosableIterator<T> recordIterator;
+ return FSUtils.isLogFile(storagePath)
+ ? createLogFileIterator(pathInfo, storagePath, fileSize)
+ : createBaseFileIterator(pathInfo, storagePath, fileSize);
+ }
+
+ /**
+ * Reads a base file using the engine's schema-evolution support, matching
+ * {@code HoodieFileGroupReader#makeBaseFileIterator}.
+ */
+ private ClosableIterator<BufferedRecord<T>>
createBaseFileIterator(StoragePathInfo pathInfo,
+
StoragePath storagePath,
+ long
fileSize) throws IOException {
+ ClosableIterator<T> recordIterator = getFileRecordIterator(
+ pathInfo,
+ storagePath,
+ fileSize,
+ readerContext.getSchemaHandler().getTableSchema(),
+ readerSchema);
+ return toBufferedRecordIterator(recordIterator, readerSchema);
+ }
+
+ /**
+ * Reads a native data log with the same schema flow as an inline log block.
+ *
+ * <p>With schema evolution enabled, records are decoded with the writer
schema stored in the
+ * native log footer and then transformed once to the evolved schema for
that instant. Without
+ * schema evolution, the file reader projects directly to the required
reader schema.
+ */
+ private ClosableIterator<BufferedRecord<T>>
createLogFileIterator(StoragePathInfo pathInfo,
+
StoragePath storagePath,
+ long
fileSize) throws IOException {
+ if (readerContext.getSchemaHandler().getInternalSchema().isEmptySchema()) {
+ ClosableIterator<T> recordIterator = getFileRecordIterator(
+ pathInfo,
+ storagePath,
+ fileSize,
+ readerContext.getSchemaHandler().getTableSchema(),
+ readerSchema);
+ return toBufferedRecordIterator(recordIterator, readerSchema);
+ }
+
+ HoodieSchema writerSchema =
TableSchemaResolver.readSchemaFromLogFile(metaClient, storagePath);
Review Comment:
add comments.
##########
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java:
##########
@@ -359,12 +363,22 @@ public static HoodieTableMetaClient initTableIfNotExists(
basePath, conf.get(FlinkOptions.TABLE_NAME));
}
- return StreamerUtil.createMetaClient(conf, hadoopConf);
+ HoodieTableMetaClient metaClient = StreamerUtil.createMetaClient(conf,
hadoopConf);
+ validateInsertOperationStorageLayout(conf,
metaClient.getTableConfig().getTableStorageLayout());
+ return metaClient;
// Do not close the filesystem in order to use the CACHE,
// some filesystems release the handles in #close method.
}
+ private static void validateInsertOperationStorageLayout(
Review Comment:
move to HoodieTableFactory.
##########
hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java:
##########
@@ -168,11 +195,20 @@ HoodieSchema generateRequiredSchema(DeleteContext
deleteContext) {
boolean hasInstantRange = readerContext.getInstantRange().isPresent();
//might need to change this if other queries than mor have mandatory fields
if (!readerContext.getHasLogFiles()) {
+ List<HoodieSchemaField> addedFields = new ArrayList<>();
if (hasInstantRange && !findNestedField(requestedSchema,
HoodieRecord.COMMIT_TIME_METADATA_FIELD).isPresent()) {
- List<HoodieSchemaField> addedFields =
Collections.singletonList(getField(this.tableSchema,
HoodieRecord.COMMIT_TIME_METADATA_FIELD));
- return appendFieldsToSchemaDedupNested(requestedSchema, addedFields);
+ addedFields.add(getField(this.tableSchema,
HoodieRecord.COMMIT_TIME_METADATA_FIELD));
}
- return requestedSchema;
+ // LSM readers merge sorted runs by key even when a file slice contains
only a base file.
Review Comment:
revert.
##########
hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java:
##########
@@ -314,32 +315,85 @@ private ClosableIterator<BufferedRecord<T>>
createFileIterator(StoragePathInfo p
if (FSUtils.isNativeDeleteLogFile(storagePath.getName())) {
return createNativeDeleteLogIterator(pathInfo, storagePath, fileSize);
}
- Pair<HoodieSchema, Map<String, String>> requiredSchemaAndRenamedColumns =
-
readerContext.getSchemaHandler().getRequiredSchemaForFileAndRenamedColumns(storagePath);
- HoodieSchema fileRequiredSchema =
requiredSchemaAndRenamedColumns.getLeft();
- ClosableIterator<T> recordIterator;
+ return FSUtils.isLogFile(storagePath)
+ ? createLogFileIterator(pathInfo, storagePath, fileSize)
+ : createBaseFileIterator(pathInfo, storagePath, fileSize);
+ }
+
+ /**
+ * Reads a base file using the engine's schema-evolution support, matching
+ * {@code HoodieFileGroupReader#makeBaseFileIterator}.
+ */
+ private ClosableIterator<BufferedRecord<T>>
createBaseFileIterator(StoragePathInfo pathInfo,
+
StoragePath storagePath,
+ long
fileSize) throws IOException {
+ ClosableIterator<T> recordIterator = getFileRecordIterator(
+ pathInfo,
+ storagePath,
+ fileSize,
+ readerContext.getSchemaHandler().getTableSchema(),
+ readerSchema);
+ return toBufferedRecordIterator(recordIterator, readerSchema);
+ }
+
+ /**
+ * Reads a native data log with the same schema flow as an inline log block.
+ *
+ * <p>With schema evolution enabled, records are decoded with the writer
schema stored in the
+ * native log footer and then transformed once to the evolved schema for
that instant. Without
+ * schema evolution, the file reader projects directly to the required
reader schema.
+ */
+ private ClosableIterator<BufferedRecord<T>>
createLogFileIterator(StoragePathInfo pathInfo,
+
StoragePath storagePath,
+ long
fileSize) throws IOException {
+ if (readerContext.getSchemaHandler().getInternalSchema().isEmptySchema()) {
+ ClosableIterator<T> recordIterator = getFileRecordIterator(
+ pathInfo,
+ storagePath,
+ fileSize,
+ readerContext.getSchemaHandler().getTableSchema(),
+ readerSchema);
+ return toBufferedRecordIterator(recordIterator, readerSchema);
+ }
+
+ HoodieSchema writerSchema =
TableSchemaResolver.readSchemaFromLogFile(metaClient, storagePath);
+ Pair<Function<T, T>, HoodieSchema> schemaEvolutionTransformer =
+ readerContext.getSchemaHandler().getSchemaEvolutionTransformer(
+ writerSchema, FSUtils.getCommitTime(storagePath.getName())).get();
+ ClosableIterator<T> recordIterator = getFileRecordIterator(
+ pathInfo, storagePath, fileSize, writerSchema, writerSchema);
+ recordIterator = new CloseableMappingIterator<>(recordIterator,
schemaEvolutionTransformer.getLeft());
+ return toBufferedRecordIterator(recordIterator,
schemaEvolutionTransformer.getRight());
+ }
+
+ private ClosableIterator<T> getFileRecordIterator(StoragePathInfo pathInfo,
+ StoragePath storagePath,
+ long fileSize,
+ HoodieSchema dataSchema,
+ HoodieSchema
requiredSchema) throws IOException {
if (pathInfo != null) {
- recordIterator = readerContext.getFileRecordIterator(
- pathInfo, 0, pathInfo.getLength(),
readerContext.getSchemaHandler().getTableSchema(), fileRequiredSchema, storage);
+ return readerContext.getFileRecordIterator(
+ pathInfo, 0, pathInfo.getLength(), dataSchema, requiredSchema,
storage);
} else {
long length = fileSize >= 0 ? fileSize :
storage.getPathInfo(storagePath).getLength();
- recordIterator = readerContext.getFileRecordIterator(
- storagePath, 0, length,
readerContext.getSchemaHandler().getTableSchema(), fileRequiredSchema, storage);
- }
- if (!areSchemasProjectionEquivalent(fileRequiredSchema, readerSchema) ||
!requiredSchemaAndRenamedColumns.getRight().isEmpty()) {
- UnaryOperator<T> projector = readerContext.getRecordContext()
- .projectRecord(fileRequiredSchema, readerSchema,
requiredSchemaAndRenamedColumns.getRight());
- recordIterator = new CloseableMappingIterator<>(recordIterator,
projector);
+ return readerContext.getFileRecordIterator(
+ storagePath, 0, length, dataSchema, requiredSchema, storage);
}
+ }
+
+ private ClosableIterator<BufferedRecord<T>>
toBufferedRecordIterator(ClosableIterator<T> recordIterator,
+
HoodieSchema recordSchema) {
if (readerContext.getInstantRange().isPresent()) {
recordIterator = readerContext.applyInstantRangeFilter(recordIterator);
}
return new CloseableMappingIterator<>(recordIterator, record ->
BufferedRecords.fromEngineRecord(
- readerContext.getRecordContext().seal(record),
- readerSchema,
+ record,
+ recordSchema,
readerContext.getRecordContext(),
orderingFieldNames,
- readerContext.getRecordContext().isDeleteRecord(record,
readerContext.getSchemaHandler().getDeleteContext().withReaderSchema(readerSchema))));
+ readerContext.getRecordContext().isDeleteRecord(
+ record,
readerContext.getSchemaHandler().getDeleteContext().withReaderSchema(recordSchema)))
+ .toBinary(readerContext.getRecordContext()));
Review Comment:
todo: avoid toBinary
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]