vishalmore90 opened a new pull request, #40285:
URL: https://github.com/apache/beam/pull/40285

   ### Context
   Fixes #40284. 
   In the Python SDK, `apache_beam.io.parquetio._RowDictionariesToArrowTable` 
(used for writing records to PyArrow tables) suffered from severe cross-window 
state leakage. Previously, elements processed within the same bundle were 
buffered into a single, flat list regardless of their respective windows or 
panes. In `finish_bundle()`, the entire aggregated PyArrow table was emitted 
into a single `WindowedValue` corresponding *only* to the window of the most 
recently processed element in that bundle, and `PaneInfo` was discarded. This 
caused data to be assigned to incorrect windows and panes, silently corrupting 
downstream aggregations in unbounded/streaming scenarios.
   
   ### Changes
   *   **`sdks/python/apache_beam/io/parquetio.py`**: Refactored 
`_RowDictionariesToArrowTable` state management.
   *   Replaced flat `_buffer`, `_record_batches`, and 
`_record_batches_byte_size` attributes with dictionaries keyed by a `(window, 
pane)` tuple.
   *   In `process()`, elements are now appended strictly to the buffer 
assigned to their specific window and pane.
   *   In `finish_bundle()`, a separate PyArrow table and `WindowedValue` is 
yielded for each `(window, pane)` key, properly preserving `pane_info` metadata.
   
   ### Verification
   *   Verified that the `DoFn` maintains identical functionality for 
bounded/single-window pipelines.
   *   Confirmed isolated state buckets correctly prevent multi-window bundles 
from overwriting previously processed windows.
   *   The `TODO` complaining about missing pane info has been inherently 
resolved, as `pane_info` is now yielded natively.
   
   ### PR Checklist
   - [x] I have read the 
[CONTRIBUTING.md](https://github.com/apache/beam/blob/master/CONTRIBUTING.md) 
and [Code of 
Conduct](https://github.com/apache/beam/blob/master/CODE_OF_CONDUCT.md).
   - [x] I have analyzed the root cause and implemented a minimal, safe fix.
   - [x] The fix addresses the exact issue without regressions or style 
violations.
   


-- 
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