vishalmore90 opened a new issue, #40284:
URL: https://github.com/apache/beam/issues/40284

   ### What happened?
   
   
   In the Python SDK, `apache_beam.io.parquetio._WriteBatches` (the core `DoFn` 
that buffers and converts records into PyArrow tables) suffers from 
cross-window state leakage when processing elements from multiple windows 
within a single bundle.
   
   The `DoFn` buffers incoming rows into a single, flat `self._buffer` 
irrespective of the `window` they belong to. When it assigns `self._window = w` 
in `process()`, it completely overwrites the window from previously processed 
elements in the same bundle. 
   
   Upon `finish_bundle()`, the entire aggregated PyArrow table (which may 
contain elements from `WindowA`, `WindowB`, etc.) is emitted inside a single 
`WindowedValue` matching only the *last* seen window (`self._window`). 
Furthermore, it completely drops the `PaneInfo`, implicitly reverting to 
default pane metadata. A `TODO` explicitly marks confusion over pane info 
retrieval (`TODO(pabloem) HOW DO WE GET THE PANE`), highlighting this 
architectural flaw.
   
   ## Code Pointers / Steps to Reproduce
   The flaw is in `sdks/python/apache_beam/io/parquetio.py` in the 
`_WriteBatches` class:
   
   ```python
     def process(self, row, w=DoFn.WindowParam, pane=DoFn.PaneInfoParam):
       # Bug: Overwrites window for the entire bundle, losing granularity
       self._window = w 
       
       # ... buffers row into self._buffer WITHOUT grouping by window ...
   
     def finish_bundle(self):
         # ...
         else:
           # unbounded input
           yield WindowedValue(
               table,
               timestamp=self._window.end, 
               # Bug: All rows in the bundle are assigned to the last seen 
window.
               windows=[self._window]  # TODO(pabloem) HOW DO WE GET THE PANE
           )
   ```
   
   ## Impact
   In streaming or unbounded scenarios where a bundle can contain elements from 
multiple windows (or sliding windows where elements belong to multiple windows 
simultaneously), downstream aggregations will be corrupted. Data belonging to 
older windows will be incorrectly injected into the latest window seen in the 
bundle, causing data loss for earlier windows and incorrect 
metrics/aggregations downstream.
   
   ## Proposed Solution
   The `_WriteBatches` `DoFn` must group buffers by window (similar to how 
`_WriteWindowedBundleDoFn` in `iobase.py` manages state using `w_key`). 
   1. `self._buffer`, `self._record_batches`, and 
`self._record_batches_byte_size` should be maintained as dictionaries keyed by 
the `Window` object.
   2. In `process()`, append elements to the buffer specifically assigned to 
`w`.
   3. In `finish_bundle()`, iterate over the dictionary and emit a separate 
`WindowedValue` for each window.
   4. `PaneInfo` must be captured and managed per-window, or correctly 
propagated when emitting the `WindowedValue`.
   
   
   ### Issue Priority
   
   Priority: 2 (default / most bugs should be filed as P2)
   
   ### Issue Components
   
   - [x] Component: Python SDK
   - [ ] Component: Java SDK
   - [ ] Component: Go SDK
   - [ ] Component: Typescript SDK
   - [ ] Component: IO connector
   - [ ] Component: Beam YAML
   - [ ] Component: Beam examples
   - [ ] Component: Beam playground
   - [ ] Component: Beam katas
   - [ ] Component: Website
   - [ ] Component: Infrastructure
   - [ ] Component: Spark Runner
   - [ ] Component: Flink Runner
   - [ ] Component: Prism Runner
   - [ ] Component: Twister2 Runner
   - [ ] Component: Hazelcast Jet Runner
   - [ ] Component: Google Cloud Dataflow Runner


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