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:
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).
self._buffer, self._record_batches, and self._record_batches_byte_size should be maintained as dictionaries keyed by the Window object.
- In
process(), append elements to the buffer specifically assigned to w.
- In
finish_bundle(), iterate over the dictionary and emit a separate WindowedValue for each window.
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
What happened?
In the Python SDK,
apache_beam.io.parquetio._WriteBatches(the coreDoFnthat buffers and converts records into PyArrow tables) suffers from cross-window state leakage when processing elements from multiple windows within a single bundle.The
DoFnbuffers incoming rows into a single, flatself._bufferirrespective of thewindowthey belong to. When it assignsself._window = winprocess(), 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 fromWindowA,WindowB, etc.) is emitted inside a singleWindowedValuematching only the last seen window (self._window). Furthermore, it completely drops thePaneInfo, implicitly reverting to default pane metadata. ATODOexplicitly 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.pyin the_WriteBatchesclass: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
_WriteBatchesDoFnmust group buffers by window (similar to how_WriteWindowedBundleDoFniniobase.pymanages state usingw_key).self._buffer,self._record_batches, andself._record_batches_byte_sizeshould be maintained as dictionaries keyed by theWindowobject.process(), append elements to the buffer specifically assigned tow.finish_bundle(), iterate over the dictionary and emit a separateWindowedValuefor each window.PaneInfomust be captured and managed per-window, or correctly propagated when emitting theWindowedValue.Issue Priority
Priority: 2 (default / most bugs should be filed as P2)
Issue Components