developer-rpai commented on issue #40284:
URL: https://github.com/apache/beam/issues/40284#issuecomment-5843772662

   I've opened PR #40293 with a fix.
   
   **Confirmed on master:** in  (iobase.py), the sink's  is applied as a  
*before* the window-grouping , so  routinely sees multiple windows per bundle 
in streaming pipelines. It buffered all rows into one flat buffer, overwrote  
per element, and emitted everything in a single  with the last window seen.
   
   **Fix:** buffer state (, , ) is now keyed by window — mirroring how  keys 
its writers — and  (plus mid-bundle row-group flushes) emits one  per window. 
Pane info is intentionally not propagated since rows from many panes can land 
in one table and the downstream file sink keys writers by window only (this 
also resolves the ). The bounded/GlobalWindow path is unchanged.
   
   **Tests:** added  with regression tests (interleaved windows, mid-bundle 
flush isolation, global-window passthrough). I also executed the edited class 
directly with real PyArrow across 5 scenarios — all pass, and the pre-fix code 
fails the interleaved-windows case as expected.
   
   Happy to adjust if reviewers prefer a different shape.


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