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]
