apache/beam · error · RuntimeError

Trying to reset a non-cleared ListBuffer.

Error message

Trying to reset a non-cleared ListBuffer.

What it means

ListBuffer.reset() is meant to revive a cleared buffer for reuse; it raises RuntimeError when called on a buffer that is not cleared. Resetting a live buffer would silently discard the lifecycle invariant (cleared flag toggling), so it is only allowed after clear().

Solutions

  1. Only call reset() on buffers you previously clear()ed; check buffer.cleared first.
  2. For new buffers, skip reset() entirely — they are already usable.
  3. Track cleared state in the pool and only reset on checkout of a cleared buffer.

Example fix

# before
buf.reset()  # fresh buffer -> RuntimeError
# after
if buf.cleared:
  buf.reset()
Defensive patterns

Strategy: type-guard

Validate before calling

if buffer.cleared:
    buffer.reset()

Type guard

def needs_reset(buf):
    return buf.cleared

Try / catch

try:
    buffer.reset()
except RuntimeError:
    pass  # buffer was never cleared; nothing to reset

Prevention

When it happens

Trigger: Calling buffer.reset() on a fresh or actively-used buffer (cleared is False), e.g. unconditionally calling reset() between bundles on a newly created buffer.

Common situations: Buffer-reuse pooling code that calls reset() defensively on every checkout, including brand-new buffers; double reset() calls after an earlier reset already flipped cleared to False.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


AI-assisted analysis of apache/beam@12126d8942 (2026-09-13). Data as JSON: /api/errors/cdf56da33a1e1f07. Report an issue: GitHub.

Appendix: source

Thrown at sdks/python/apache_beam/runners/portability/fn_api_runner/execution.py:181

            idx = (idx + 1) % n
        self._grouped_output = [[output_stream.get()]
                                for output_stream in output_stream_list]
      return self._grouped_output

  def __iter__(self) -> Iterator[bytes]:
    if self.cleared:
      raise RuntimeError('Trying to iterate through a cleared ListBuffer.')
    return iter(self._inputs)

  def clear(self) -> None:
    self.cleared = True
    self._inputs = []
    self._grouped_output = None

  def reset(self) -> None:
    """Resets a cleared buffer for reuse."""
    if not self.cleared:
      raise RuntimeError('Trying to reset a non-cleared ListBuffer.')
    self.cleared = False


class GroupingBuffer(object):
  """Used to accumulate groupded (shuffled) results."""
  def __init__(
      self,
      pre_grouped_coder: coders.Coder,
      post_grouped_coder: coders.Coder,
      windowing: core.Windowing) -> None:
    self._key_coder = pre_grouped_coder.key_coder()
    self._pre_grouped_coder = pre_grouped_coder
    self._post_grouped_coder = post_grouped_coder
    self._table: collections.defaultdict[bytes,
                                         list[Any]] = collections.defaultdict(
                                             list)
    self._windowing = windowing
    self._grouped_output: Optional[list[list[bytes]]] = None

View on GitHub (pinned to 12126d8942)