Skip to content
Open
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 @@ -1000,6 +1000,40 @@ public long estimateDataSizeByListingFiles(ConnectorSession session, ConnectorTa
}
}

@Override
public long estimateDataSizeByListingFiles(ConnectorSession session, ConnectorTableHandle handle,
List<String> selectedPartitionNames) {
if (!(handle instanceof HiveTableHandle)) {
return siblingMetadata(session, handle)
.estimateDataSizeByListingFiles(session, handle, selectedPartitionNames);
}
HiveTableHandle hiveHandle = (HiveTableHandle) handle;
if (hiveHandle.getTableType() != HiveTableType.HIVE) {
return -1;
}
if (selectedPartitionNames.isEmpty()) {
return 0;
}
if (hiveHandle.isTransactional()) {
return -1;
}

ClassLoader previous = Thread.currentThread().getContextClassLoader();
try {
Thread.currentThread().setContextClassLoader(getClass().getClassLoader());
FileSystem fs = storage().getFileSystem(session);
return estimateSelectedPartitionsDataSize(hiveHandle, selectedPartitionNames,
STATS_PARTITION_SAMPLE_SIZE,
(location, values) -> sumCachedFileSizes(hiveHandle, location, values, fs));
} catch (RuntimeException e) {
LOG.warn("Failed to estimate selected hive partition data size for {}.{} from file list",
hiveHandle.getDbName(), hiveHandle.getTableName(), e);
return -1;
} finally {
Thread.currentThread().setContextClassLoader(previous);
}
}

/**
* Returns the raw byte length of every data file across ALL partitions (not sampled, not summed), a port of
* legacy {@code HMSExternalTable.getChunkSizes} for {@code ANALYZE ... WITH SAMPLE}. Only plain-hive tables
Expand Down Expand Up @@ -1079,6 +1113,49 @@ long estimateDataSize(HiveTableHandle handle, int sampleSize, ToLongBiFunction<S
}
}

long estimateSelectedPartitionsDataSize(HiveTableHandle handle, List<String> selectedPartitionNames,
int sampleSize, ToLongBiFunction<String, List<String>> sizeOf) {
try {
if (selectedPartitionNames.isEmpty()) {
return 0;
}
int selectedPartitionCount = selectedPartitionNames.size();
boolean sampled = sampleSize > 0 && sampleSize < selectedPartitionCount;
List<String> chosenPartitionNames = selectedPartitionNames;
if (sampled) {
List<String> shuffled = new ArrayList<>(selectedPartitionNames);
Collections.shuffle(shuffled);
chosenPartitionNames = shuffled.subList(0, sampleSize);
}
List<HmsPartitionInfo> partitions = hmsClient.getExistingPartitions(
handle.getDbName(), handle.getTableName(), chosenPartitionNames);
if (partitions.size() != chosenPartitionNames.size()) {
return -1;
}
List<PartitionRef> refs = new ArrayList<>(partitions.size());
for (HmsPartitionInfo partition : partitions) {
String location = partition.getLocation();
if (location == null || location.isEmpty()) {
return -1;
}
refs.add(new PartitionRef(location, partition.getValues()));
}

long selectedSize = 0;
for (PartitionRef ref : refs) {
selectedSize += Math.max(0, sizeOf.applyAsLong(ref.location, ref.partitionValues));
}
if (sampled && selectedSize == 0) {
return -1;
}
return sampled ? scaleSampledSize(selectedSize, selectedPartitionCount, refs.size()) : selectedSize;
} catch (RuntimeException e) {
LOG.warn("Failed to estimate selected hive partition data size for {}.{} from file list",
handle.getDbName(), handle.getTableName(), e);
return -1;
}
}

/**
* Scales a sampled data size up to the whole table: {@code sampledSize * totalPartitions /
* sampledPartitions} (legacy {@code HMSExternalTable.getRowCountFromFileList}). Multiplies BEFORE dividing
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import org.apache.doris.connector.hms.HmsDatabaseInfo;
import org.apache.doris.connector.hms.HmsPartitionInfo;
import org.apache.doris.connector.hms.HmsTableInfo;
import org.apache.doris.connector.spi.ConnectorSession;
import org.apache.doris.filesystem.FileSystem;

import org.junit.jupiter.api.Assertions;
Expand Down Expand Up @@ -95,6 +96,68 @@ public void sampledPartitionsAreScaledUpToTheWholeTable() {
Assertions.assertEquals(400L, metadata(client).estimateDataSize(partitioned(), 2, (loc, vals) -> 100));
}

@Test
public void selectedPartitionEstimateOnlyListsTheSelectedRange() {
PartitionFakeHmsClient client = new PartitionFakeHmsClient(Arrays.asList("p0", "p1", "p2"));
long size = metadata(client).estimateSelectedPartitionsDataSize(
partitioned(), Collections.singletonList("p1"), 30,
(loc, vals) -> loc.endsWith("p1") ? 200 : 10_000);

Assertions.assertEquals(200L, size);
}

@Test
public void selectedPartitionSamplingScalesWithinTheSelectedRange() {
PartitionFakeHmsClient client = new PartitionFakeHmsClient(
Arrays.asList("p0", "p1", "p2", "p3", "p4"));
long size = metadata(client).estimateSelectedPartitionsDataSize(
partitioned(), Arrays.asList("p1", "p2", "p3", "p4"), 2,
(loc, vals) -> 100);

Assertions.assertEquals(400L, size);
Assertions.assertEquals(2, client.lastRequestedPartitionNames.size());
Assertions.assertTrue(Arrays.asList("p1", "p2", "p3", "p4")
.containsAll(client.lastRequestedPartitionNames));
}

@Test
public void zeroSizedSampleDoesNotProveTheWholeSelectedRangeIsEmpty() {
PartitionFakeHmsClient client = new PartitionFakeHmsClient(
Arrays.asList("p0", "p1", "p2", "p3"));
long size = metadata(client).estimateSelectedPartitionsDataSize(
partitioned(), Arrays.asList("p0", "p1", "p2", "p3"), 2,
(loc, vals) -> 0);

Assertions.assertEquals(-1L, size);
}

@Test
public void fullyInspectedEmptySelectedPartitionsHaveZeroDataSize() {
PartitionFakeHmsClient client = new PartitionFakeHmsClient(Arrays.asList("p0", "p1"));
long size = metadata(client).estimateSelectedPartitionsDataSize(
partitioned(), Arrays.asList("p0", "p1"), 30,
(loc, vals) -> 0);

Assertions.assertEquals(0L, size);
}

@Test
public void emptySelectedPartitionRangeHasZeroDataSize() {
PartitionFakeHmsClient client = new PartitionFakeHmsClient(Collections.emptyList());
Assertions.assertEquals(0L, metadata(client).estimateSelectedPartitionsDataSize(
partitioned(), Collections.emptyList(), 30, (loc, vals) -> 100));
Assertions.assertTrue(client.lastRequestedPartitionNames.isEmpty());
}

@Test
public void missingSelectedPartitionMakesTheEstimateUnknown() {
PartitionFakeHmsClient client = new PartitionFakeHmsClient(Arrays.asList("p0", "p1"));
client.removeExistingPartition("p1");

Assertions.assertEquals(-1L, metadata(client).estimateSelectedPartitionsDataSize(
partitioned(), Arrays.asList("p0", "p1"), 30, (loc, vals) -> 100));
}

@Test
public void zeroTotalSizeReturnsMinusOne() {
PartitionFakeHmsClient client = new PartitionFakeHmsClient(Arrays.asList("p0", "p1"));
Expand Down Expand Up @@ -193,6 +256,27 @@ public void nonHiveTableTypeIsNotEstimated() {
.estimateDataSizeByListingFiles(null, hudiHandle));
}

@Test
public void transactionalHiveSelectedPartitionsAreNotEstimated() {
HiveTableHandle transactionalHandle = new HiveTableHandle.Builder("db", "t", HiveTableType.HIVE)
.partitionKeyNames(Collections.singletonList("dt"))
.tableParameters(Collections.singletonMap("transactional", "true"))
.build();
PartitionFakeHmsClient client = new PartitionFakeHmsClient(Collections.singletonList("p0"));
FakeConnectorContext context = new FakeConnectorContext() {
@Override
public FileSystem getFileSystem(ConnectorSession session) {
return new FakeFileSystem();
}
};
HiveConnectorMetadata metadata = new HiveConnectorMetadata(
client, HiveTestProperties.minimal(), context);

Assertions.assertEquals(-1L, metadata.estimateDataSizeByListingFiles(
null, transactionalHandle, Collections.singletonList("p0")));
Assertions.assertTrue(client.lastRequestedPartitionNames.isEmpty());
}

/** A {@link HiveFileListingCache} whose listing always fails, to prove listFileSizes propagates (not swallows). */
private static final class ThrowingFileListingCache extends HiveFileListingCache {
ThrowingFileListingCache() {
Expand All @@ -214,6 +298,8 @@ public List<HiveFileStatus> listDataFiles(String dbName, String tableName, Strin
private static final class PartitionFakeHmsClient implements HmsClient {
private final List<String> partitionNames;
private final java.util.Set<String> withoutLocation = new java.util.HashSet<>();
private final java.util.Set<String> missingPartitions = new java.util.HashSet<>();
private List<String> lastRequestedPartitionNames = Collections.emptyList();

PartitionFakeHmsClient(List<String> partitionNames) {
this.partitionNames = partitionNames;
Expand All @@ -223,6 +309,10 @@ void dropLocationFor(String name) {
withoutLocation.add(name);
}

void removeExistingPartition(String name) {
missingPartitions.add(name);
}

@Override
public List<String> listPartitionNames(String dbName, String tableName, int maxParts) {
return partitionNames;
Expand All @@ -239,6 +329,19 @@ public List<HmsPartitionInfo> getPartitions(String dbName, String tableName, Lis
return result;
}

@Override
public List<HmsPartitionInfo> getExistingPartitions(
String dbName, String tableName, List<String> partNames) {
lastRequestedPartitionNames = new ArrayList<>(partNames);
List<String> existingNames = new ArrayList<>();
for (String name : partNames) {
if (!missingPartitions.contains(name)) {
existingNames.add(name);
}
}
return getPartitions(dbName, tableName, existingNames);
}

@Override
public HmsTableInfo getTable(String dbName, String tableName) {
throw new UnsupportedOperationException();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -241,6 +241,8 @@ public void everyPerHandleMethodForwardsAForeignHandleToTheSibling() {
md.getTableStatistics(session, foreignHandle);
md.getColumnStatistics(session, foreignHandle, "c");
long size = md.estimateDataSizeByListingFiles(session, foreignHandle);
long selectedSize = md.estimateDataSizeByListingFiles(
session, foreignHandle, Collections.singletonList("p"));
Optional<FilterApplicationResult<ConnectorTableHandle>> filter = md.applyFilter(session, foreignHandle, null);
List<String> partNames = md.listPartitionNames(session, foreignHandle);
md.listPartitions(session, foreignHandle, Optional.empty());
Expand Down Expand Up @@ -278,6 +280,9 @@ public void everyPerHandleMethodForwardsAForeignHandleToTheSibling() {
// A few return values prove the ANSWER is the sibling's, not hive's default.
Assertions.assertEquals(RecordingSiblingMetadata.SENTINEL_SIZE, size,
"estimateDataSize must return the sibling's value, not hive's -1");
Assertions.assertEquals(RecordingSiblingMetadata.SENTINEL_SIZE, selectedSize,
"selected-partition estimateDataSize must return the sibling's value, not hive's -1");
Assertions.assertEquals(Collections.singletonList("p"), siblingMetadata.selectedPartitionNames);
Assertions.assertEquals(RecordingSiblingMetadata.SENTINEL_SNAPSHOT_ID, pin.getSnapshotId(),
"beginQuerySnapshot must return the sibling's snapshot-id pin, not hive's -1 last-modified pin");
Assertions.assertEquals(Collections.singletonMap("p", 55L), partitionFreshness,
Expand Down Expand Up @@ -726,7 +731,7 @@ private static final class RecordingSiblingMetadata implements ConnectorMetadata
// dropping a guard, or adding one that should not forward, changes this list and fails the test).
static final List<String> EXPECTED_METHODS = Collections.unmodifiableList(Arrays.asList(
"getTableSchema", "getColumnHandles", "getTableStatistics", "getColumnStatistics",
"estimateDataSizeByListingFiles",
"estimateDataSizeByListingFiles", "estimateSelectedDataSizeByListingFiles",
"applyFilter", "listPartitionNames", "listPartitions",
"beginQuerySnapshot", "getTableFreshness", "getPartitionFreshnessMillis",
"getPartitionsFreshnessMillis", "dropTable",
Expand All @@ -747,6 +752,7 @@ private static final class RecordingSiblingMetadata implements ConnectorMetadata
"validateRowLevelDmlMode", "validateStaticPartitionColumns", "validateWritePartitionNames"));

final List<String> calls = new ArrayList<>();
List<String> selectedPartitionNames = Collections.emptyList();
final Optional<FilterApplicationResult<ConnectorTableHandle>> filterResult =
Optional.of(new FilterApplicationResult<>(SIBLING_HANDLE, null, false));

Expand Down Expand Up @@ -792,6 +798,14 @@ public long estimateDataSizeByListingFiles(ConnectorSession session, ConnectorTa
return SENTINEL_SIZE;
}

@Override
public long estimateDataSizeByListingFiles(ConnectorSession session, ConnectorTableHandle handle,
List<String> selectedPartitionNames) {
calls.add("estimateSelectedDataSizeByListingFiles");
this.selectedPartitionNames = new ArrayList<>(selectedPartitionNames);
return SENTINEL_SIZE;
}


@Override
public Optional<FilterApplicationResult<ConnectorTableHandle>> applyFilter(ConnectorSession session,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1724,7 +1724,7 @@ public void planWriteThreadsPinnedReadSnapshotFromHandleToTransaction() {
}

@Test
public void planMergePreservesExplicitlyEmptyReadAcrossConcurrentFirstAppend() {
public void planMergeKeepsExplicitEmptyReadFencedAcrossConcurrentFirstAppend() {
InMemoryCatalog catalog = freshCatalog();
TableIdentifier id = TableIdentifier.of("db1", "tv2");
Table empty = catalog.createTable(id, SCHEMA, PartitionSpec.unpartitioned(),
Expand All @@ -1751,8 +1751,9 @@ public void planMergePreservesExplicitlyEmptyReadAcrossConcurrentFirstAppend() {
providerFor(ops.table, ctx).planWrite(new WriteSession(txn),
new WriteHandle(emptyPinnedHandle).writeOperation(WriteOperation.MERGE));

Assertions.assertNull(txn.getBaseSnapshotId(),
"an explicitly empty read must leave RowDelta validation unbounded across the first append");
Assertions.assertEquals(Long.valueOf(-1L), txn.getBaseSnapshotId(),
"an explicit empty read is an OCC fence (base -1), not an absent pin: the pinned empty "
+ "generation must survive the concurrent first append instead of drifting to S1");
}

// ───────────────────────────── MERGE sink (TIcebergMergeSink) ─────────────────────────────
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,23 @@ default long estimateDataSizeByListingFiles(ConnectorSession session, ConnectorT
return -1;
}

/**
* Estimates the on-disk data size in bytes for exactly the named partitions. The names are connector
* partition identifiers returned to the engine; an empty list means that pruning selected no partitions.
* Returns -1 when the connector cannot estimate this exact partition set. Zero is a valid result and means
* the selected partitions contain no data files. The supplied handle already carries any snapshot pin for
* this scan; an implementation that cannot honor that exact scope must return -1.
*
* <p>The default must not delegate to the whole-table overload because that would silently substitute an
* incompatible row-count scope.</p>
*/
default long estimateDataSizeByListingFiles(
ConnectorSession session,
ConnectorTableHandle handle,
List<String> selectedPartitionNames) {
return -1;
}

/**
* Returns the RAW byte length of every data file across ALL partitions of the table (not sampled, not summed),
* for {@code ANALYZE ... WITH SAMPLE}: fe-core seed-shuffles and cumulates these sizes to a sample scale
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,9 +87,9 @@ public void connectorApiMajorTracksTheRecordedSurfaceChange() throws IOException
Assertions.assertNotNull(in, "missing connector plugin API version resource");
version.load(in);
}
// Major 12 adds the SUPPORTS_FIELD_ID_ACCESS_PATH and SUPPORTS_SYS_TABLE_NESTED_COLUMN_PRUNE
// capabilities: a plugin naming either constant cannot link against an older FE.
Assertions.assertEquals("12.0", version.getProperty("api.version"));
// Major 13 在 major 12 的接口基础上新增所选分区的数据量估算重载。
// 旧 FE 必须在链接不兼容字节码前拒绝使用该重载的插件。
Assertions.assertEquals("13.0", version.getProperty("api.version"));
}

/** Root entry points plus provider/handle types returned to connector plugins. */
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ dropTable(org.apache.doris.connector.spi.ConnectorSession,org.apache.doris.conne
dropTag(org.apache.doris.connector.spi.ConnectorSession,org.apache.doris.connector.spi.handle.ConnectorTableHandle,org.apache.doris.connector.spi.ddl.DropRefChange)
dropView(org.apache.doris.connector.spi.ConnectorSession,java.lang.String,java.lang.String)
estimateDataSizeByListingFiles(org.apache.doris.connector.spi.ConnectorSession,org.apache.doris.connector.spi.handle.ConnectorTableHandle)
estimateDataSizeByListingFiles(org.apache.doris.connector.spi.ConnectorSession,org.apache.doris.connector.spi.handle.ConnectorTableHandle,java.util.List)
fromRemoteColumnName(org.apache.doris.connector.spi.ConnectorSession,java.lang.String,java.lang.String,java.lang.String)
fromRemoteDatabaseName(org.apache.doris.connector.spi.ConnectorSession,java.lang.String)
fromRemoteTableName(org.apache.doris.connector.spi.ConnectorSession,java.lang.String,java.lang.String)
Expand Down
2 changes: 1 addition & 1 deletion fe/fe-connector/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ under the License.
of the latter two means bumping this property as well (and fe-extension-spi means bumping
all five families).
-->
<connector.plugin.api.version>12.0</connector.plugin.api.version>
<connector.plugin.api.version>13.0</connector.plugin.api.version>
</properties>

<modules>
Expand Down
Loading
Loading