developer-rpai opened a new pull request, #40293:
URL: https://github.com/apache/beam/pull/40293

   ## Problem
   
   In the Python SDK, `_RowDictionariesToArrowTable` 
(`sdks/python/apache_beam/io/parquetio.py`) — the `DoFn` that buffers rows and 
converts them into PyArrow tables on the `WriteToParquet` path — suffers 
cross-window state leakage when a bundle contains elements from multiple 
windows.
   
   It buffered incoming rows into a single flat `self._buffer` irrespective of 
window, overwrote `self._window` on every `process()` call, and at 
`finish_bundle()` emitted the entire mixed-window table inside one 
`WindowedValue` carrying only the **last** window seen (plus a `TODO(pabloem) 
HOW DO WE GET THE PANE`).
   
   This is reachable in practice: in `WriteImpl._expand_unbounded` 
(`iobase.py`), the sink's `convert_fn` is applied as a `ParDo` **before** the 
window-grouping `GroupByKey`, so a single bundle routinely spans windows in 
streaming pipelines. Downstream, the mislabeled table is grouped into the wrong 
window's shard: data loss for earlier windows and silently corrupted 
aggregations.
   
   Fixes #40284.
   
   ## Fix
   
   - Buffer state (`_buffers`, `_record_batches`, `_record_batches_byte_size`) 
is now keyed by window, mirroring how `_WriteWindowedBundleDoFn` keys its 
writers by window.
   - `process()` appends rows to the current element's window buffer; 
mid-bundle row-group flushes yield a `WindowedValue` for that window.
   - `finish_bundle()` emits one `WindowedValue` per window and resets state; 
`start_bundle()` also resets for DoFn reuse safety.
   - Pane info is intentionally not propagated: rows from many panes may be 
batched into one table, and the downstream file sink keys its writers by window 
only. The confusing `TODO(pabloem)` is resolved by this design.
   - The bounded/GlobalWindow path is unchanged (single windowed value at end 
of window).
   
   ## Tests
   
   Added `RowDictionariesToArrowTableTest` to `parquetio_test.py` with three 
regression tests:
   - interleaved windows in one bundle produce one `WindowedValue` per window 
with correctly attributed rows (fails on the old code: 1 output instead of 2);
   - mid-bundle buffer flushes stay within their window;
   - the GlobalWindow path still emits a single windowed value.
   
   ## Verification level
   
   I verified the fix by executing the real edited class (extracted verbatim 
from the file) with real PyArrow in 5 scenarios: interleaved windows, 
mid-bundle buffer flush, global window, empty bundle, and byte-size-triggered 
mid-bundle flush — all pass, and the pre-fix code fails the interleaved-windows 
scenario as expected. I could not run the repo's full `parquetio_test.py` suite 
here (no Beam SDK install in this environment); CI should run it.
   


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