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]
