Skip to content

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

Description

@vishalmore90

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

  1. self._buffer, self._record_batches, and self._record_batches_byte_size should be maintained as dictionaries keyed by the Window object.
  2. In process(), append elements to the buffer specifically assigned to w.
  3. In finish_bundle(), iterate over the dictionary and emit a separate WindowedValue for each window.
  4. 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

  • Component: Python SDK
  • Component: Java SDK
  • Component: Go SDK
  • Component: Typescript SDK
  • Component: IO connector
  • Component: Beam YAML
  • Component: Beam examples
  • Component: Beam playground
  • Component: Beam katas
  • Component: Website
  • Component: Infrastructure
  • Component: Spark Runner
  • Component: Flink Runner
  • Component: Prism Runner
  • Component: Twister2 Runner
  • Component: Hazelcast Jet Runner
  • Component: Google Cloud Dataflow Runner

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions