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]