Skip to content
Draft
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
Original file line number Diff line number Diff 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 number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
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.RecordWriter;
Expand All @@ -36,13 +37,17 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
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.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 Down Expand Up @@ -122,8 +127,35 @@ public CommitMessage doCompact(FileStoreTable table, String commitUser) throws E
dataFileMeta =
dataFileMeta.assignSequenceNumber(
minSequenceId(compactBefore), maxSequenceId(compactBefore));
dataFileMeta =
dataFileMeta.withColumnMaxSequenceNumbers(
compactedColumnMaxSequenceNumbers(table, dataFileMeta));
compactAfter.add(dataFileMeta);

return commitMessage(compactBefore, compactAfter);
}

private long[] compactedColumnMaxSequenceNumbers(
FileStoreTable table, DataFileMeta outputFile) {
Map<Integer, Long> fieldMaxSequences = new HashMap<>();
for (DataFileMeta input : compactBefore) {
List<DataField> inputFields = fileFields(table.schemaManager()::schema, input);
for (int inputPosition = 0; inputPosition < inputFields.size(); inputPosition++) {
fieldMaxSequences.merge(
inputFields.get(inputPosition).id(),
fieldMaxSequenceNumber(input, inputPosition, inputFields.size()),
(left, right) -> Math.max(left, right));
}
}

long fallbackSequence = maxSequenceId(compactBefore);
List<DataField> outputFields = fileFields(table.schemaManager()::schema, outputFile);
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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,8 @@
import java.util.Set;
import java.util.TreeMap;

import static org.apache.paimon.utils.DataEvolutionUtils.fileFieldIds;
import static org.apache.paimon.utils.DataEvolutionUtils.fieldMaxSequenceNumber;
import static org.apache.paimon.utils.DataEvolutionUtils.fileFields;

/** Plans existing global index files which need refresh after data-evolution updates. */
public final class DataEvolutionGlobalIndexRefreshPlanner {
Expand Down Expand Up @@ -81,7 +82,7 @@ public static List<IndexManifestEntry> findIndexesToRefresh(
.addIndex(i, indexMeta.rowRange(), scanSnapshotId);
}

Map<Pair<Long, List<String>>, Set<Integer>> fileFieldIdsCache = new HashMap<>();
Map<Pair<Long, List<String>>, List<DataField>> fileFieldsCache = new HashMap<>();
for (ManifestEntry dataEntry : dataEntries) {
DataFileMeta file = dataEntry.file();
if (dataEntry.kind() != FileKind.ADD || file.firstRowId() == null) {
Expand All @@ -93,12 +94,21 @@ public static List<IndexManifestEntry> findIndexesToRefresh(
continue;
}

Set<Integer> physicalFieldIds =
fileFieldIdsCache.computeIfAbsent(
List<DataField> physicalFields =
fileFieldsCache.computeIfAbsent(
Pair.of(file.schemaId(), file.writeCols()),
key -> fileFieldIds(schemaManager::schema, file));
if (!disjoint(indexedFieldIds, physicalFieldIds)) {
group.addDataFile(file);
key -> fileFields(schemaManager::schema, file));
long indexedMaxSequence = Long.MIN_VALUE;
for (int position = 0; position < physicalFields.size(); position++) {
if (indexedFieldIds.contains(physicalFields.get(position).id())) {
indexedMaxSequence =
Math.max(
indexedMaxSequence,
fieldMaxSequenceNumber(file, position, physicalFields.size()));
}
}
if (indexedMaxSequence != Long.MIN_VALUE) {
group.addDataFile(file, indexedMaxSequence);
}
}

Expand All @@ -119,7 +129,7 @@ public static List<IndexManifestEntry> findIndexesToRefresh(
private static final class RefreshGroup {

private final List<IndexQuery> indexes = new ArrayList<>();
private final List<DataFileMeta> dataFiles = new ArrayList<>();
private final List<DataUpdate> dataUpdates = new ArrayList<>();
private final MergedRanges indexedRanges = new MergedRanges();
private long minScanSnapshotId = Long.MAX_VALUE;

Expand All @@ -134,21 +144,23 @@ private boolean mayContainUpdate(DataFileMeta file) {
&& indexedRanges.intersects(file.nonNullRowIdRange());
}

private void addDataFile(DataFileMeta file) {
dataFiles.add(file);
private void addDataFile(DataFileMeta file, long maxSequenceNumber) {
if (maxSequenceNumber > minScanSnapshotId) {
dataUpdates.add(new DataUpdate(file.nonNullRowIdRange(), maxSequenceNumber));
}
}

private void markIndexesToRefresh(boolean[] result) {
// As scan watermarks decrease, eligible data files only grow.
dataFiles.sort(Comparator.comparingLong(DataFileMeta::maxSequenceNumber).reversed());
dataUpdates.sort(Comparator.comparingLong(DataUpdate::maxSequenceNumber).reversed());
indexes.sort((left, right) -> Long.compare(right.scanSnapshotId, left.scanSnapshotId));

MergedRanges updatedRanges = new MergedRanges();
int nextFile = 0;
for (IndexQuery index : indexes) {
while (nextFile < dataFiles.size()
&& dataFiles.get(nextFile).maxSequenceNumber() > index.scanSnapshotId) {
updatedRanges.add(dataFiles.get(nextFile).nonNullRowIdRange());
while (nextFile < dataUpdates.size()
&& dataUpdates.get(nextFile).maxSequenceNumber > index.scanSnapshotId) {
updatedRanges.add(dataUpdates.get(nextFile).rowRange);
nextFile++;
}
if (updatedRanges.intersects(index.rowRange)) {
Expand All @@ -158,6 +170,21 @@ private void markIndexesToRefresh(boolean[] result) {
}
}

private static final class DataUpdate {

private final Range rowRange;
private final long maxSequenceNumber;

private DataUpdate(Range rowRange, long maxSequenceNumber) {
this.rowRange = rowRange;
this.maxSequenceNumber = maxSequenceNumber;
}

private long maxSequenceNumber() {
return maxSequenceNumber;
}
}

private static final class IndexQuery {

private final int ordinal;
Expand Down Expand Up @@ -218,13 +245,4 @@ private static boolean matchesFields(GlobalIndexMeta meta, List<DataField> field
}
return expectedExtraFields != null && Arrays.equals(actualExtraFields, expectedExtraFields);
}

private static boolean disjoint(Set<Integer> left, Set<Integer> right) {
for (Integer value : left) {
if (right.contains(value)) {
return false;
}
}
return true;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,22 @@ public List<String> writeCols() {
return nullableStringArray(Fields.WRITE_COLS);
}

@Nullable
@Override
public long[] columnMaxSequenceNumbers() {
int position = requiredPosition(Fields.COLUMN_MAX_SEQUENCE_NUMBERS);
InternalRow row = currentRow();
if (row.isNullAt(position)) {
return null;
}
InternalArray array = row.getArray(position);
long[] result = new long[array.size()];
for (int i = 0; i < array.size(); i++) {
result[i] = array.getLong(i);
}
return result;
}

public boolean containsWriteColumn(BinaryString fieldName) {
int position = requiredPosition(Fields.WRITE_COLS);
InternalRow row = currentRow();
Expand Down Expand Up @@ -281,6 +297,11 @@ public DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenc
throw unsupportedOperation("assignSequenceNumber(long, long)");
}

@Override
public DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) {
throw unsupportedOperation("withColumnMaxSequenceNumbers(long[])");
}

@Override
public DataFileMeta assignFirstRowId(long firstRowId) {
throw unsupportedOperation("assignFirstRowId(long)");
Expand Down Expand Up @@ -378,6 +399,8 @@ private static class Fields {
private static final int EXTERNAL_PATH = fieldIndex(DataFileMeta.EXTERNAL_PATH);
private static final int FIRST_ROW_ID = fieldIndex(DataFileMeta.FIRST_ROW_ID);
private static final int WRITE_COLS = fieldIndex(DataFileMeta.WRITE_COLS);
private static final int COLUMN_MAX_SEQUENCE_NUMBERS =
fieldIndex(DataFileMeta.COLUMN_MAX_SEQUENCE_NUMBERS);
}

/** Projected data-file schema together with its bound binary field layout. */
Expand Down
75 changes: 72 additions & 3 deletions paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
Original file line number Diff line number Diff 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 COLUMN_MAX_SEQUENCE_NUMBERS = "_COLUMN_MAX_SEQUENCE_NUMBERS";

RowType SCHEMA =
new RowType(
Expand Down Expand 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,
COLUMN_MAX_SEQUENCE_NUMBERS,
new ArrayType(true, new BigIntType(false)))));

BinaryRow EMPTY_MIN_KEY = EMPTY_ROW;
BinaryRow EMPTY_MAX_KEY = EMPTY_ROW;
Expand Down Expand Up @@ -173,7 +178,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 +198,54 @@ static DataFileMeta create(
valueStatsCols,
externalPath,
firstRowId,
writeCols);
writeCols,
null);
}

static DataFileMeta create(
String fileName,
long fileSize,
long rowCount,
BinaryRow minKey,
BinaryRow maxKey,
SimpleStats keyStats,
SimpleStats valueStats,
long minSequenceNumber,
long maxSequenceNumber,
long schemaId,
int level,
List<String> extraFiles,
Timestamp creationTime,
@Nullable Long deleteRowCount,
@Nullable byte[] embeddedIndex,
@Nullable FileSource fileSource,
@Nullable List<String> valueStatsCols,
@Nullable String externalPath,
@Nullable Long firstRowId,
@Nullable List<String> writeCols,
@Nullable long[] columnMaxSequenceNumbers) {
return new PojoDataFileMeta(
fileName,
fileSize,
rowCount,
minKey,
maxKey,
keyStats,
valueStats,
minSequenceNumber,
maxSequenceNumber,
schemaId,
level,
extraFiles,
creationTime,
deleteRowCount,
embeddedIndex,
fileSource,
valueStatsCols,
externalPath,
firstRowId,
writeCols,
columnMaxSequenceNumbers);
}

static DataFileMeta create(
Expand Down Expand Up @@ -354,6 +406,18 @@ 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
default long[] columnMaxSequenceNumbers() {
return null;
}

DataFileMeta upgrade(int newLevel);

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

DataFileMeta assignSequenceNumber(long minSequenceNumber, long maxSequenceNumber);

default DataFileMeta withColumnMaxSequenceNumbers(long[] columnMaxSequenceNumbers) {
throw new UnsupportedOperationException(
"This DataFileMeta implementation does not support column sequence numbers.");
}

DataFileMeta assignFirstRowId(long firstRowId);

DataFileMeta newFirstRowId(@Nullable Long newFirstRowId);
Expand Down
Original file line number Diff line number Diff 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 Down
Loading
Loading