Skip to content

[Python] Fix ParquetIO silently corrupting data across windows in unbounded pipelines - #40293

Open
developer-rpai wants to merge 4 commits into
apache:masterfrom
developer-rpai:fix-40284-parquetio-windowed-buffers
Open

developer-rpai wants to merge 4 commits into
apache:masterfrom
developer-rpai:fix-40284-parquetio-windowed-buffers

Conversation

@developer-rpai

Copy link
Copy Markdown

Problem

In the Python SDK, _RowDictionariesToArrowTable (sdks/python/apache_beam/io/parquetio.py) — the DoFn that buffers rows and converts them into PyArrow tables on the WriteToParquet path — suffers cross-window state leakage when a bundle contains elements from multiple windows.

It buffered incoming rows into a single flat self._buffer irrespective of window, overwrote self._window on every process() call, and at finish_bundle() emitted the entire mixed-window table inside one WindowedValue carrying only the last window seen (plus a TODO(pabloem) HOW DO WE GET THE PANE).

This is reachable in practice: in WriteImpl._expand_unbounded (iobase.py), the sink's convert_fn is applied as a ParDo before the window-grouping GroupByKey, 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

  • Buffer state (_buffers, _record_batches, _record_batches_byte_size) is now keyed by window, mirroring how _WriteWindowedBundleDoFn keys its writers by window.
  • process() appends rows to the current element's window buffer; mid-bundle row-group flushes yield a WindowedValue for that window.
  • finish_bundle() emits one WindowedValue per window and resets state; start_bundle() also resets for DoFn reuse safety.
  • Pane info is intentionally not propagated: rows from many panes may be batched into one table, and the downstream file sink keys its writers by window only. The confusing TODO(pabloem) is resolved by this design.
  • The bounded/GlobalWindow path is unchanged (single windowed value at end of window).

Tests

Added RowDictionariesToArrowTableTest to parquetio_test.py with three regression tests:

  • interleaved windows in one bundle produce one WindowedValue per window with correctly attributed rows (fails on the old code: 1 output instead of 2);
  • mid-bundle buffer flushes stay within their window;
  • the GlobalWindow path still emits a single windowed value.

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.py suite here (no Beam SDK install in this environment); CI should run it.

_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
@github-actions

Copy link
Copy Markdown
Contributor

Assigning reviewers:

R: @claudevdm for label python.

Note: If you would like to opt out of this review, comment assign to next reviewer.

Available commands:

  • stop reviewer notifications - opt out of the automated review tooling
  • remind me after tests pass - tag the comment author after tests pass
  • waiting on author - shift the attention set back to the author (any comment or push by the author will return the attention set to the reviewers)

The PR bot will only process comments in the main thread (not review comments).

@claudevdm claudevdm left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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)
@developer-rpai

Copy link
Copy Markdown
Author

Thanks for the thorough review, @claudevdm. Good catch on the sink -- fixed as suggested: open() now returns (writer, file_handle) and the handle travels with the writer instead of living on the sink, so each window's writer closes its own file and we no longer get Parquet files with missing footers.

Also added the end-to-end test_write_streaming_rows_land_in_their_own_window test to WriteStreamingTest (one bundle spanning 3 windows, each row verified in its own window's file). The new RowDictionariesToArrowTableTest unit tests cover the DoFn-level windowing behavior.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Bug]: Python ParquetIO _WriteBatches DoFn silently corrupts data across windows in unbounded pipelines

2 participants