fdolce opened a new pull request, #29381:
URL: https://github.com/apache/flink/pull/29381
## What is the purpose of the change
Adds ways to consume DataFrame results incrementally or partially, instead
of only via `collect()` / `to_pandas()`, which materialize the whole result and
never return on unbounded sources. Users can now stream rows or batches from a
running job, or peek at the first `n` rows with an optional timeout, and the
job is cancelled when they're done.
## Brief change log
- `DataFrame.iter_rows()` / `DataFrame.iter_batches()` return a
`CloseableIterator` of dicts or pandas/PyArrow batches; closing it (or leaving
the `with` block) cancels the job
- `DataFrame.take(n)` / `DataFrame.take_batch(n)` return the first `n`
rows and cancel the job; an optional `timeout` returns whatever arrived by then
- All four can add each row's change kind (`+I`/`-U`/`+U`/`-D`) via
`include_row_kind`, for updating results
- Batch schemas come from the Flink schema, including nested
ROW/ARRAY/MAP; unsupported types are rejected before a job is submitted
- New `pyflink.dataframe.validation` module for shared argument checks
Also changed here:
- `limit` / `offset` / `head` now use the shared integer check; the error
for a negative `n` changed from "n must be non-negative" to "n must be at least
0". Let me know if we prefer to keep this out of this PR though
### Known limitation: TIMESTAMP_LTZ
The new methods read results through `TableResult.collect()`, the same path
as `DataFrame.collect()`. That path currently cannot transfer `TIMESTAMP_LTZ`
values: `PythonBridgeUtils.getPickledBytesFromRow` fails to pickle them. This
already affects `collect()` in both the DataFrame and Table APIs, and this PR
doesn't change it. `to_pandas()` is unaffected because it uses the Arrow path.
For now this is documented in the docstrings of the four new methods rather
than checked up front, so the error appears once the job runs. It's common in
practice because `from_records` types Python `datetime` values as
`TIMESTAMP_LTZ`. Until it's fixed, cast the column to `TIMESTAMP` or `STRING`.
I'd prefer to have a general fix for the conversion in a separate PR rather
than add a DataFrame-only check here.
## Verifying this change
This change added tests and can be verified as follows:
- `pyflink/dataframe/tests/test_iteration.py` covers:
- rows/batches match `to_pandas()`, batch sizing, empty batches, nested
types
- job cancellation on early close, on `take` completion, and on timeout
against an unbounded datagen source
- job failures are re-raised and the job terminates
- change kinds of an updating aggregation
- invalid arguments (including non-finite timeouts) are rejected before
any job is submitted
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency): no
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: yes (new `@PublicEvolving` methods on `DataFrame` and new
`CloseableIterator`)
- The serializers: no
- The runtime per-record code paths (performance sensitive): no
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
- The S3 file system connector: no
## Documentation
- Does this pull request introduce a new feature? yes
- If yes, how is the feature documented? docs (PyDoc + API reference in
`dataframe.rst`)
---
##### Was generative AI tooling used to co-author this PR?
- [X] Yes (please specify the tool below)
Generated-by: Claude Code (Claude Opus 5.5)
--
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]