Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -3734,6 +3734,10 @@ public GlobalIndexColumnUpdateAction globalIndexColumnUpdateAction() {
return options.get(GLOBAL_INDEX_COLUMN_UPDATE_ACTION);
}

public boolean ignoreIndexColumnUpdate() {
return globalIndexColumnUpdateAction() == GlobalIndexColumnUpdateAction.IGNORE;
}

public LookupStrategy lookupStrategy() {
return LookupStrategy.from(
mergeEngine().equals(MergeEngine.FIRST_ROW),
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -40,7 +40,7 @@
public class DataEvolutionCompactTaskSerializer
implements VersionedSerializer<DataEvolutionCompactTask> {

private static final int CURRENT_VERSION = 2;
private static final int CURRENT_VERSION = 3;

private final DataFileMetaSerializer dataFileSerializer;

Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,25 +26,36 @@
import org.apache.paimon.manifest.FileSource;
import org.apache.paimon.operation.AppendFileStoreWrite;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.SchemaManager;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.sink.CommitMessage;
import org.apache.paimon.table.source.DataSplit;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.FileStorePathFactory;
import org.apache.paimon.utils.Pair;
import org.apache.paimon.utils.RecordWriter;
import org.apache.paimon.utils.SetUtils;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import javax.annotation.Nullable;

import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Function;
import java.util.stream.Collectors;

import static org.apache.paimon.types.BlobType.fieldNamesInBlobFile;
import static org.apache.paimon.types.VectorType.fieldNamesInVectorFile;
import static org.apache.paimon.types.VectorType.isVectorStoreFile;
import static org.apache.paimon.utils.DataEvolutionUtils.checkContiguousRowRange;
import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber;
import static org.apache.paimon.utils.DataEvolutionUtils.fileFields;
import static org.apache.paimon.utils.Preconditions.checkArgument;

/** Compacts normal structured files of a data evolution table. */
Expand DownExpand Up@@ -126,8 +137,64 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E
dataFileMeta =
dataFileMeta.assignSequenceNumber(
minSequenceId(compactBefore), maxSequenceId(compactBefore));
if (options.ignoreIndexColumnUpdate()) {
long[] columnMaxSequenceNumbers =
compactedColumnMaxSequenceNumbers(table, dataFileMeta);
if (columnMaxSequenceNumbers != null) {
dataFileMeta = dataFileMeta.withColumnMaxSequenceNumbers(columnMaxSequenceNumbers);
}
}
compactAfter.add(dataFileMeta);

return commitMessage(compactBefore, compactAfter);
}

@Nullable
private long[] compactedColumnMaxSequenceNumbers(
FileStoreTable table, DataFileMeta outputFile) {
SchemaManager schemaManager = table.schemaManager();
Map<Long, TableSchema> schemaCache = new HashMap<>();
Function<Long, TableSchema> schemaLoader =
schemaId -> schemaCache.computeIfAbsent(schemaId, schemaManager::schema);
Map<Pair<Long, List<String>>, List<DataField>> fileFieldsCache = new HashMap<>();

Map<Integer, Long> fieldMaxSequences = new HashMap<>();
for (DataFileMeta input : compactBefore) {
List<DataField> inputFields =
fileFieldsCache.computeIfAbsent(
Pair.of(input.schemaId(), input.writeCols()),
key -> fileFields(schemaLoader, input));
long[] inputColumnSequences = input.columnMaxSequenceNumbers();
for (int inputPosition = 0; inputPosition < inputFields.size(); inputPosition++) {
fieldMaxSequences.merge(
inputFields.get(inputPosition).id(),
fieldMaxSequenceNumber(
input, inputColumnSequences, inputPosition, inputFields.size()),
Math::max);
}
}

long fallbackSequence = outputFile.maxSequenceNumber();
List<DataField> outputFields =
fileFieldsCache.computeIfAbsent(
Pair.of(outputFile.schemaId(), outputFile.writeCols()),
key -> fileFields(schemaLoader, outputFile));
boolean allEqualToFileMax =
outputFields.stream()
.allMatch(
field ->
fieldMaxSequences.getOrDefault(field.id(), fallbackSequence)
== fallbackSequence);
if (allEqualToFileMax) {
return null;
}

long[] result = new long[outputFields.size()];
for (int outputPosition = 0; outputPosition < outputFields.size(); outputPosition++) {
result[outputPosition] =
fieldMaxSequences.getOrDefault(
outputFields.get(outputPosition).id(), fallbackSequence);
}
return result;
}
}
36 changes: 29 additions & 7 deletions paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
Original file line numberDiff line numberDiff line change
Expand Up@@ -80,6 +80,7 @@ public interface DataFileMeta {
String EXTERNAL_PATH = "_EXTERNAL_PATH";
String FIRST_ROW_ID = "_FIRST_ROW_ID";
String WRITE_COLS = "_WRITE_COLS";
String WRITE_COLS_SEQUENCES = "_WRITE_COLS_SEQUENCES";

RowType SCHEMA =
new RowType(
Expand DownExpand Up@@ -109,7 +110,11 @@ public interface DataFileMeta {
new DataField(17, EXTERNAL_PATH, newStringType(true)),
new DataField(18, FIRST_ROW_ID, new BigIntType(true)),
new DataField(
19, WRITE_COLS, new ArrayType(true, newStringType(false)))));
19, WRITE_COLS, new ArrayType(true, newStringType(false))),
new DataField(
20,
WRITE_COLS_SEQUENCES,
new ArrayType(true, new BigIntType(false)))));

BinaryRow EMPTY_MIN_KEY = EMPTY_ROW;
BinaryRow EMPTY_MAX_KEY = EMPTY_ROW;
Expand DownExpand Up@@ -150,7 +155,8 @@ static DataFileMeta forAppend(
valueStatsCols,
externalPath,
firstRowId,
writeCols);
writeCols,
null);
}

static DataFileMeta create(
Expand All@@ -173,7 +179,7 @@ static DataFileMeta create(
@Nullable String externalPath,
@Nullable Long firstRowId,
@Nullable List<String> writeCols) {
return new PojoDataFileMeta(
return create(
fileName,
fileSize,
rowCount,
Expand All@@ -193,7 +199,8 @@ static DataFileMeta create(
valueStatsCols,
externalPath,
firstRowId,
writeCols);
writeCols,
null);
}

static DataFileMeta create(
Expand DownExpand Up@@ -234,7 +241,8 @@ static DataFileMeta create(
valueStatsCols,
null,
firstRowId,
writeCols);
writeCols,
null);
}

static DataFileMeta create(
Expand All@@ -257,7 +265,8 @@ static DataFileMeta create(
@Nullable List<String> valueStatsCols,
@Nullable String externalPath,
@Nullable Long firstRowId,
@Nullable List<String> writeCols) {
@Nullable List<String> writeCols,
@Nullable long[] columnMaxSequenceNumbers) {
return new PojoDataFileMeta(
fileName,
fileSize,
Expand All@@ -278,7 +287,8 @@ static DataFileMeta create(
valueStatsCols,
externalPath,
firstRowId,
writeCols);
writeCols,
columnMaxSequenceNumbers);
}

String fileName();
Expand DownExpand Up@@ -354,6 +364,16 @@ default Range nonNullRowIdRange() {
@Nullable
List<String> writeCols();

/**
* Maximum sequence number per physical table field after data-evolution compaction.
*
* <p>Values follow the table-field order selected by {@link #writeCols()} when it is non-null
* (system fields are ignored), or the file schema field order otherwise. A null value means
* that only the file-level sequence range is available.
*/
@Nullable
long[] columnMaxSequenceNumbers();

DataFileMeta upgrade(int newLevel);

DataFileMeta rename(String newFileName);
Expand All@@ -362,6 +382,8 @@ default Range nonNullRowIdRange() {

DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenceNumber);

DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers);

DataFileMeta assignFirstRowId(long firstRowId);

DataFileMeta newFirstRowId(@Nullable Long newFirstRowId);
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -136,6 +136,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException {
null,
null,
null,
null,
null);
}
}
Original file line numberDiff line numberDiff line change
Expand Up@@ -142,6 +142,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException {
null,
null,
null,
null,
null);
}
}
Original file line numberDiff line numberDiff line change
Expand Up@@ -147,6 +147,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException {
row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)),
null,
null,
null,
null);
}
}
Original file line numberDiff line numberDiff line change
Expand Up@@ -149,6 +149,7 @@ public DataFileMeta deserialize(DataInputView in) throws IOException {
row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)),
row.isNullAt(17) ? null : row.getString(17).toString(),
null,
null,
null);
}
}
Original file line numberDiff line numberDiff line change
Expand Up@@ -36,7 +36,7 @@ public class DataFileMetaFirstRowIdLegacySerializer extends ObjectSerializer<Dat
private static final long serialVersionUID = 1L;

public DataFileMetaFirstRowIdLegacySerializer() {
super(DataFileMeta.SCHEMA);
super(DataFileMetaWriteColsLegacySerializer.SCHEMA);
}

@Override
Expand DownExpand Up@@ -86,6 +86,7 @@ public DataFileMeta fromRow(InternalRow row) {
row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)),
row.isNullAt(17) ? null : row.getString(17).toString(),
row.isNullAt(18) ? null : row.getLong(18),
null,
null);
}
}
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,7 +19,9 @@
package org.apache.paimon.io;

import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalArray;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.manifest.FileSource;
import org.apache.paimon.stats.SimpleStats;
Expand All@@ -41,6 +43,7 @@ public DataFileMetaSerializer() {

@Override
public InternalRow toRow(DataFileMeta meta) {
long[] columnMaxSequenceNumbers = meta.columnMaxSequenceNumbers();
return GenericRow.of(
BinaryString.fromString(meta.fileName()),
meta.fileSize(),
Expand All@@ -61,7 +64,10 @@ public InternalRow toRow(DataFileMeta meta) {
toStringArrayData(meta.valueStatsCols()),
meta.externalPath().map(BinaryString::fromString).orElse(null),
meta.firstRowId(),
meta.writeCols() == null ? null : toStringArrayData(meta.writeCols()));
meta.writeCols() == null ? null : toStringArrayData(meta.writeCols()),
columnMaxSequenceNumbers == null
? null
: new GenericArray(columnMaxSequenceNumbers));
}

@Override
Expand All@@ -86,6 +92,15 @@ public DataFileMeta fromRow(InternalRow row) {
row.isNullAt(16) ? null : fromStringArrayData(row.getArray(16)),
row.isNullAt(17) ? null : row.getString(17).toString(),
row.isNullAt(18) ? null : row.getLong(18),
row.isNullAt(19) ? null : fromStringArrayData(row.getArray(19)));
row.isNullAt(19) ? null : fromStringArrayData(row.getArray(19)),
row.isNullAt(20) ? null : fromLongArray(row.getArray(20)));
}

private static long[] fromLongArray(InternalArray array) {
long[] result = new long[array.size()];
for (int i = 0; i < array.size(); i++) {
result[i] = array.getLong(i);
}
return result;
}
}
Loading
Loading