claudevdm opened a new pull request, #40252:
URL: https://github.com/apache/beam/pull/40252
When a file is registered without `location_prefix`, `AddFiles` infers its
partition from the column bounds (min/max) in the Parquet footer. Today that
inference, and the bounds themselves, go wrong in several cases: some files are
registered under the wrong partition with no error, some fail the whole
pipeline, and some are refused for no good reason. The bounds bug also affects
unpartitioned tables, because query planners use the stored bounds to skip
files.
This PR:
- stores column bounds in the unit of the table column's type (millis and
nanos timestamps and times, unsigned 32-bit ints, v3 `timestamp_ns` columns);
- puts a file in the null partition only when its partition column is proven
null, and refuses (error row) files whose partition can't be determined;
- sets the inferred partition as values instead of round-tripping it through
a path string.
Wrong identity partitions are worse than they look: Iceberg's readers take
an identity-partitioned column's value from the partition tuple, not from the
file. On master, a file written without statistics with `flag = [true, true]`
is registered under `flag=false`, and every reader returns `[false, false]`.
## What changes for users
Partition inference (partitioned tables, no `location_prefix`):
| Case | master | this PR |
|---|---|---|
| Partition column all null | boolean -> `false`, string -> `"null"`,
date/int -> refused (`Text 'null' could not be parsed`) | null partition |
| File lacks the partition column (written before it existed) | boolean ->
`false` | null partition |
| Partition column without statistics (writer stats off, INT96) | treated as
all null, as above | error row: no column bounds, use a location prefix |
| Partition column with both nulls and values | registered under the value's
partition: `IS NULL` misses the null rows, `= value` returns them | error row |
| Avro file | `NullPointerException` fails the pipeline | error row |
| `year(ts)` | stored as ordinal 2024 (year 3994) | correct |
| `month`, `hour` or `identity` on a timestamp | refused (`For input string:
"2024-01"`) | registered |
| Table metrics mode `none`, file has a uint32 or millis `time` column |
`ClassCastException` fails the pipeline | works |
Column bounds (all tables):
| File column | master | this PR |
|---|---|---|
| Timestamp in millis or nanos, time in nanos | raw file unit, off by 1000x:
range filters can skip the file, `day(ts)` is wrong | table column's unit |
| Time in millis (INT32), unsigned INT32 | `ClassCastException`, file
refused | bounds computed (under a `time` / `long` column) |
| Millis or micros timestamp under a v3 `timestamp_ns` column | 10^6x /
1000x too small | nanos |
The first commit only adds tests (all but one control fail on master); each
later commit makes a group of them pass and names them in its message.
1. **Regression tests** (new `AddFilesMetricsTest`). Bounds tests call
`getFileMetrics` directly. Partition tests run `AddFiles` in a `TestPipeline`
and check the error output and the committed `DataFile`'s partition. Test files
have no field IDs, like Spark or pyarrow output, so they resolve through the
name mapping.
2. **Never guess a partition for a column without bounds.** Missing bounds
used to mean "all null". Now the null partition needs proof: an empty file, a
null count equal to the value count (`allValuesNull`), or a Parquet file
without the column (`lacksColumn`; such a column reads as null). Anything else
is an unknown partition. Bound maps are read null-safely for Avro (`orEmpty`).
Worth checking: the null partition is chosen only when every row is null.
3. **Set the partition as values.** `getPartitionFromMetrics` returns the
`PartitionKey`, and the `DataFile` is built with `withPartition` instead of
`withPartitionPath(pk.toPath())`. The text round-trip wrote null as `null`
(parsed back as `false` or the string `"null"`) and couldn't parse `month`,
`hour` or timestamp values. The `location_prefix` path is unchanged.
`SerializableDataFile` already carries the partition as JSON, so the null
reaches the commit.
4. **Second metrics pass for partition columns only.** When the table's
metrics mode drops a partition column's bounds, AddFiles computes them from the
same footer for inference only (they aren't stored). That pass now defaults
every other column to `none`, so a column Iceberg can't handle no longer fails
the bundle.
5. **Bounds in the table column's unit** (`BoundAdjustment`). Iceberg's
converter types every Parquet timestamp and time as micros and uint32 as long,
but copies the raw statistic. The conversion is chosen from the file's
annotation and the table column's type, which `getFileMetrics` now receives:
x1000, /1000 (lower rounded down, upper up), x10^6, or unsigned. INT32 columns
that would throw (millis `time`, uint32) are shown to Iceberg without their
annotation so it computes int bounds, then converted. A bound that overflows
the table's unit (e.g. year 9999 in nanos) is dropped. If the table column
isn't the matching type (e.g. a timestamp file column under a `long`), the
bounds stay as Iceberg computes them, as on master. Worth checking: the
annotation stripping, the rounding direction and the overflow handling.
6. **Nulls and values.** Bounds ignore nulls, so a partition column holding
both is now an unknown partition. The `void` transform is exempt: it maps
everything to null.
------------------------
Thank you for your contribution! Follow this checklist to help us
incorporate your contribution quickly and easily:
- [ ] Mention the appropriate issue in your description (for example:
`addresses #123`), if applicable. This will automatically add a link to the
pull request in the issue. If you would like the issue to automatically close
on merging the pull request, comment `fixes #<ISSUE NUMBER>` instead.
- [ ] Update `CHANGES.md` with noteworthy changes.
- [ ] If this contribution is large, please file an Apache [Individual
Contributor License Agreement](https://www.apache.org/licenses/icla.pdf).
See the [Contributor Guide](https://beam.apache.org/contribute) for more
tips on [how to make review process
smoother](https://github.com/apache/beam/blob/master/CONTRIBUTING.md#make-the-reviewers-job-easier).
To check the build health, please visit
[https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md](https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md)
GitHub Actions Tests Status (on master branch)
------------------------------------------------------------------------------------------------
[](https://github.com/apache/beam/actions?query=workflow%3A%22Build+python+source+distribution+and+wheels%22+branch%3Amaster+event%3Aschedule)
[](https://github.com/apache/beam/actions?query=workflow%3A%22Python+Tests%22+branch%3Amaster+event%3Aschedule)
[](https://github.com/apache/beam/actions?query=workflow%3A%22Java+Tests%22+branch%3Amaster+event%3Aschedule)
[](https://github.com/apache/beam/actions?query=workflow%3A%22Go+tests%22+branch%3Amaster+event%3Aschedule)
See [CI.md](https://github.com/apache/beam/blob/master/CI.md) for more
information about GitHub Actions CI or the [workflows
README](https://github.com/apache/beam/blob/master/.github/workflows/README.md)
to see a list of phrases to trigger workflows.
--
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]