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]

Reply via email to