nastra commented on code in PR #17433:
URL: https://github.com/apache/iceberg/pull/17433#discussion_r3749921794
##########
core/src/main/java/org/apache/iceberg/V4ManifestReader.java:
##########
@@ -292,17 +333,65 @@ private Schema readSchema(boolean hasPartitionFilter) {
if (columns != null) {
Schema selected =
caseSensitive ? fullSchema.select(columns) :
fullSchema.caseInsensitiveSelect(columns);
- return addRequiredColumns(selected, hasPartitionFilter);
+ return addRequiredColumns(fullSchema, selected, requiredFieldIds,
hasPartitionFilter);
}
if (requestedProjection != null) {
- return addRequiredColumns(requestedProjection, hasPartitionFilter);
+ return addRequiredColumns(
+ fullSchema, requestedProjection, requiredFieldIds,
hasPartitionFilter);
}
return fullSchema;
}
- private Schema addRequiredColumns(Schema projection, boolean
hasPartitionFilter) {
+ /** Returns the schema of everything this reader may read, including
content stats. */
+ private Schema fullSchema(Set<Integer> requiredStatsProjectionFieldIds) {
+ Types.StructType contentStatsType =
contentStatsType(requiredStatsProjectionFieldIds);
+ Schema base = TrackedFile.schema(unionPartitionType, contentStatsType);
+ if (contentStatsType.fields().isEmpty()) {
+ // schema uses the unknown type for empty stats, which cannot be
paired with the stats
+ // struct in the manifest, so drop the field instead of reading it as
unknown
+ base = TypeUtil.selectNot(base,
ImmutableSet.of(TrackedFile.CONTENT_STATS_ID));
+ }
+
+ // the read schema carries row_position (via BASE_TYPE) so the reader
can fill manifestPos
+ return TypeUtil.replaceFieldTypes(
+ base, ImmutableMap.of(TrackedFile.TRACKING.fieldId(),
TrackingStruct.BASE_TYPE));
+ }
+
+ /** Returns the stats type to read, which is empty when no stats are
needed. */
+ private Types.StructType contentStatsType(Set<Integer>
requiredStatsProjectionForFieldIds) {
+ if (scanPlanning || statsProjectionForFieldIds != null) {
+ // scan planning and projectStats(fieldIds) both narrow the set of
stats that are read
+ return StatsUtil.statsReadSchema(tableSchema,
requiredStatsProjectionForFieldIds);
+ }
+
+ return StatsUtil.statsReadSchema(
+ tableSchema, TypeUtil.indexById(tableSchema.asStruct()).keySet());
Review Comment:
> This walks the full table schema for every build() — TypeUtil.indexById
once, then statsReadSchema walks again (plus indexParents and per-field
isScalar climbs to the root). Fine per manifest, but this is on the default
path (no projectStats, no forScanPlanning) taken for every manifest read that
copies entries forward, so the cost multiplies across a scan's fan-out on wide
tables.
we should be able to get this down by using
`StatsUtil.statsReadSchema(tableSchema, tableSchema.idToName().keySet())` at
the very least. As a follow-up we could maybe introduce a custom visitor, which
would then hopefully result in only a single schema walk.
> manifests only store stats for a capped prefix of columns (default ~100
via MetricsConfig), so on a table with e.g. 5,000 columns the default "read all
stats" builds a stats schema with ~5,000 slots and registers 5,000
FieldStatsStruct custom types — but ~4,900 of them resolve to null at decode
time because the manifest never stored them. We're paying construction cost for
stats we know aren't there.
I don't have a good solution for this either atm, but it's a good discussion
point to talk about. Let me do some exploration on this
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]