Core: Read content stats from v4 Manifest - #17433
Conversation
c62305e to
6cf0fd0
Compare
anoopj
left a comment
There was a problem hiding this comment.
Is the file pruning coming in a followup?
4975ce5 to
33d52b2
Compare
| * fields referenced by the {@link #filter(Expression) filter} are always read. | ||
| */ | ||
| Builder projectStats(Iterable<Integer> fieldIds) { | ||
| Preconditions.checkArgument(fieldIds != null, "Invalid stats projection for field IDs: null"); |
There was a problem hiding this comment.
I think we have 4 modes relevant to stats so far, the default/CDC, scan planing, select column by name and project schema and we conditionally add required column depends on the filter. I am wondering if we want to add coverage for
- empty projectStats and filter (ok to use precondition to check if this combination does not make sense)
- valid projectStats and filter with default mode ( I think scan planning is already covered in
projectStatsAndFilterStatsAreCombined)
33d52b2 to
b8612ba
Compare
|
|
||
| /** 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) { |
There was a problem hiding this comment.
Is the asymmetry intended? A caller doing join/aggregate pushdown has to explicitly narrow via projectStats(...), but stats for row-filter-referenced columns come along for free — without any opt-in from the caller — because the default (non-scan-planning, no projectStats) reads full stats and the filter only forces its refs to be included on top of a narrower projection.
The check itself (statsProjectionForFieldIds != null) is right — requiredStatsProjectionForFieldIds may be non-empty due to the row filter, but we only want to narrow when the caller explicitly asks via projectStats(fieldIds).
There was a problem hiding this comment.
yes this is intended, but it's worth discussing whether the default should be to read all stats vs no stats
| } | ||
|
|
||
| return StatsUtil.statsReadSchema( | ||
| tableSchema, TypeUtil.indexById(tableSchema.asStruct()).keySet()); |
There was a problem hiding this 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.
Compounding this: 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.
But I don't have a good solution. Neither option below is clean:
- Using current
MetricsConfig.metricsFieldIds()at read time is per-table, not per-manifest. If the cap narrowed since the manifest was written, we silently drop stats the manifest actually holds — no correctness impact (InclusiveMetricsEvaluatortreats absent stats as "may match"), but pruning gets coarser on copy-forward and scans open more files at query time. If it widened, we still over-ask for the extra columns and get the same null-resolution waste. Not a sound signal either way. - The only truthful source is the manifest's own
content_statsschema. But peeking at that before configuring the projection means either an extra file open per manifest (drop belowInternalDatatoAvro.read/Parquet.readfor a header/footer peek, then reopen viaInternalDatawith the intersection), or extendingInternalData.ReadBuilderwith afileSchema()accessor so the projection can be picked after the header is read. Both cost something.
Flagging this to see if we can explore good alternatives — not blocking this PR.
And orthogonally, I am also wondering if we should caching the full stats read schema keyed off the Schema (like Schema.lazyIdToField)?
There was a problem hiding this comment.
Idea worth exploring: version MetricsConfig in TableMetadata alongside schemas and partition specs, and stamp each manifest file (root or leaf) with the metrics-config-id in effect when the manifest was written. Readers resolve the id and recover the exact write-time stats field IDs, and build the read schema off the intersection with tableSchema. Correct across config drift; no extra file peek at read time.
Trade-offs: spec change (metrics-configs list + current-metrics-config-id on TableMetadata, MetricsConfigParser, metrics-config-id on manifest file entries); MetricsConfig shifts from a runtime property derivation to a first-class immutable artifact; retention/dedup matters since property tweaks churn more than schema/spec changes. With metrics-config-id, schema-id, sort-id, reader should be able to faithfully reconstruct the writer MetricsConfig and content stats write schema. This also requires move the metrics config from table properties to properly versioned struct in table metadata.
Alternatives — field-IDs list in each manifest's header, or on the manifest-list entry — either force a pre-read I/O per manifest or bloat the manifest list on wide tables, so the versioned direction seems the cleanest long-term shape.
There was a problem hiding this 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
There was a problem hiding this comment.
@rdblue clarified the expected behavior, which makes sense to me.
There are two read modes
- caller explicitly project column stats (for filter, join etc.). This is already covered by the PR.
- caller just want to read all existing column stats (for manifest update carryover). Currently, we construct the content stats from all table schema fields. Ryan was suggesting that we should just use the current MetricsConfig to construct the read schema for content_stats. We only want to carry over stats based on the latest metrics config.
We probably should rename the two APIs from StatsUtil to clarify their purpose.
statsWriteSchema->currentStatsSchemastatsReadSchema->statsProjectionSchema// we might be able to leave this unchanged too.
b8612ba to
e3d6c97
Compare
V4ManifestReadernow reads thecontent_statscolumn, so eachTrackedFilecomesback with per-column bounds and counts.
The stats schema is derived from the table schema, so
builder()takes it as a secondargument. Stats for every column are read by default, which is what copying entries into
a new manifest needs. Callers that want less can narrow stats to be read via
projectStats(...).Reading stats surfaced two bugs in the copy path, which were:
ContentStatsStruct.copy()threw an NPE when a projected column had no stats in themanifest
FieldStatsStructcopiedStructLikebounds by reference. Geometry and geography bounda bounding-box struct, so under
reuseContainers()every entry reported the last row'sbounds. Bounds are now deep-copied through
StructLikeUtil.copy, andStructCopyisSerializableso a copied bound still survives serialization.Used Claude for the initial prototyping but reviewed and adjusted the code manually