[Python] Fix ParquetIO silently corrupting data across windows in unbounded pipelines - #40293
developer-rpai wants to merge 4 commits into
Conversation
_RowDictionariesToArrowTable buffered rows from every window into a single flat buffer and overwrote self._window on each process() call. At finish_bundle() the whole mixed-window table was emitted in one WindowedValue carrying only the last window seen. In unbounded pipelines a bundle routinely spans windows (the convert_fn ParDo runs before the window-grouping GroupByKey in WriteImpl._expand_unbounded), so rows from earlier windows were silently attributed to the latest window: data loss for the earlier windows and corrupted downstream aggregations. Buffer record batches per window and emit one WindowedValue per window at finish_bundle() (and on mid-bundle row-group flushes), mirroring how _WriteWindowedBundleDoFn keys its writers by window. Pane info is intentionally not propagated: rows from many panes may land in one table, and the downstream file sink keys writers by window only. Adds regression tests driving the DoFn directly with interleaved windows. Fixes apache#40284
|
Assigning reviewers: R: @claudevdm for label python. Note: If you would like to opt out of this review, comment Available commands:
The PR bot will only process comments in the main thread (not review comments). |
claudevdm
left a comment
There was a problem hiding this comment.
Thanks for tracking this down. The per-window buffering in _RowDictionariesToArrowTable looks right, and I confirmed the new unit tests fail on master and pass on this branch. At the pipeline level, though, this change exposes a second bug further down the write path. When one bundle spans several windows, the output is still wrong, just in a different way.
Repro. The new tests call the DoFn directly. This end-to-end test could go in WriteStreamingTest. It reads the output files back and checks that each row ended up in its own window's file:
def test_write_streaming_rows_land_in_their_own_window(self):
# One add_elements call is delivered as one bundle spanning 3 windows.
base = datetime(2021, 3, 1, tzinfo=pytz.UTC).timestamp()
offsets = [1, 11, 21, 3, 13, 23]
stream = TestStream().add_elements([
beam.window.TimestampedValue({'offset': o}, base + o) for o in offsets
]).advance_watermark_to_infinity()
with TestPipeline() as p:
_ = (
p
| stream
| beam.WindowInto(beam.window.FixedWindows(10))
| beam.io.WriteToParquet(
file_path_prefix=self.tempdir + '/out',
file_name_suffix='.parquet',
num_shards=1,
schema=pa.schema([('offset', pa.int64())])))
pattern = re.compile(r'.*-\[(?P<start>[\d\.]+), [\d\.]+\)-.*\.parquet$')
rows_by_window = {}
for file_name in glob.glob(self.tempdir + '/out*'):
start = float(pattern.match(file_name).group('start')) - base
rows = pq.read_table(file_name).column('offset').to_pylist()
rows_by_window[start] = sorted(rows)
self.assertEqual({0: [1, 3], 10: [11, 13], 20: [21, 23]}, rows_by_window)
Prism and BundleBasedDirectRunner give the same results:
master: {20.0: [1, 3, 11, 13, 21, 23]}
this PR: the [20, 30) file isn't valid Parquet (ArrowInvalid: ... Parquet magic bytes not found in footer), yet the pipeline still reports success.
I think we need something like this code below to ensure the right file handle is closed on the sink
--- a/sdks/python/apache_beam/io/parquetio.py
+++ b/sdks/python/apache_beam/io/parquetio.py
@@ -848,31 +848,35 @@ class _ParquetSink(filebasedsink.FileBasedSink):
"pyarrow version >= 4.x, please use a different pyarrow version. "
f"Your pyarrow version: {pa.__version__}")
self._use_compliant_nested_type = use_compliant_nested_type
- self._file_handle = None
def open(self, temp_path):
- self._file_handle = super().open(temp_path)
+ # Several writers may be open on this sink at once (one per window in
+ # streaming writes), so the file handle travels with its writer rather
+ # than being stored on the sink.
+ file_handle = super().open(temp_path)
if ARROW_MAJOR_VERSION < 4:
- return pq.ParquetWriter(
- self._file_handle,
+ writer = pq.ParquetWriter(
+ file_handle,
self._schema,
compression=self._codec,
use_deprecated_int96_timestamps=self._use_deprecated_int96_timestamps)
- return pq.ParquetWriter(
- self._file_handle,
- self._schema,
- compression=self._codec,
- use_deprecated_int96_timestamps=self._use_deprecated_int96_timestamps,
- use_compliant_nested_type=self._use_compliant_nested_type)
-
- def write_record(self, writer, table: paTable):
+ else:
+ writer = pq.ParquetWriter(
+ file_handle,
+ self._schema,
+ compression=self._codec,
+ use_deprecated_int96_timestamps=self._use_deprecated_int96_timestamps,
+ use_compliant_nested_type=self._use_compliant_nested_type)
+ return writer, file_handle
+
+ def write_record(self, writer_and_handle, table: paTable):
+ writer, _ = writer_and_handle
writer.write_table(table)
- def close(self, writer):
+ def close(self, writer_and_handle):
+ writer, file_handle = writer_and_handle
writer.close()
- if self._file_handle:
- self._file_handle.close()
- self._file_handle = None
+ file_handle.close()
…test
The _RowDictionariesToArrowTable fix exposed a second bug further down the
write path: _ParquetSink stored the open file handle on the sink
(self._file_handle), but several writers may be open on one sink at once
(one per window in streaming writes). The wrong handle got closed, producing
Parquet files with missing footers ("Parquet magic bytes not found").
Make the file handle travel with its writer: open() now returns a
(writer, file_handle) tuple, and write_record()/close() unpack it.
Also adds the end-to-end regression test from review: one bundle spanning
3 windows must land each row in its own window's file.
Co-authored-by: claudevdm (review feedback)
|
Thanks for the thorough review, @claudevdm. Good catch on the sink -- fixed as suggested: Also added the end-to-end |
Problem
In the Python SDK,
_RowDictionariesToArrowTable(sdks/python/apache_beam/io/parquetio.py) — theDoFnthat buffers rows and converts them into PyArrow tables on theWriteToParquetpath — suffers cross-window state leakage when a bundle contains elements from multiple windows.It buffered incoming rows into a single flat
self._bufferirrespective of window, overwroteself._windowon everyprocess()call, and atfinish_bundle()emitted the entire mixed-window table inside oneWindowedValuecarrying only the last window seen (plus aTODO(pabloem) HOW DO WE GET THE PANE).This is reachable in practice: in
WriteImpl._expand_unbounded(iobase.py), the sink'sconvert_fnis applied as aParDobefore the window-groupingGroupByKey, so a single bundle routinely spans windows in streaming pipelines. Downstream, the mislabeled table is grouped into the wrong window's shard: data loss for earlier windows and silently corrupted aggregations.Fixes #40284.
Fix
_buffers,_record_batches,_record_batches_byte_size) is now keyed by window, mirroring how_WriteWindowedBundleDoFnkeys its writers by window.process()appends rows to the current element's window buffer; mid-bundle row-group flushes yield aWindowedValuefor that window.finish_bundle()emits oneWindowedValueper window and resets state;start_bundle()also resets for DoFn reuse safety.TODO(pabloem)is resolved by this design.Tests
Added
RowDictionariesToArrowTableTesttoparquetio_test.pywith three regression tests:WindowedValueper window with correctly attributed rows (fails on the old code: 1 output instead of 2);Verification level
I verified the fix by executing the real edited class (extracted verbatim from the file) with real PyArrow in 5 scenarios: interleaved windows, mid-bundle buffer flush, global window, empty bundle, and byte-size-triggered mid-bundle flush — all pass, and the pre-fix code fails the interleaved-windows scenario as expected. I could not run the repo's full
parquetio_test.pysuite here (no Beam SDK install in this environment); CI should run it.