apache/beam · error · NotImplementedError

do not process elements.

Error message

%s do not process elements.

What it means

_TransformEvaluator (and subclasses that inherit its base) only handles timer firings; the base class explicitly refuses to process data elements. Calling process_element on such an evaluator always raises NotImplementedError naming the evaluator type.

Solutions

  1. Override process_element in your evaluator subclass if it should consume elements.
  2. Don't connect an input PCollection to a transform that only processes timers.
  3. Use a different evaluator/transform appropriate for element-wise processing (e.g. a ParDo evaluator).

Example fix

# before
class MyEvaluator(_TransformEvaluator):
    pass
# after
class MyEvaluator(_TransformEvaluator):
    def process_element(self, element):
        ...  # handle the element
Defensive patterns

Strategy: type-guard

Type guard

def evaluator_accepts_elements(evaluator) -> bool:
    return type(evaluator).process_element is not _TransformEvaluator.process_element

Try / catch

try:
    evaluator.process_element(element)
except NotImplementedError as e:
    # route to a timer-only evaluator or fix the wiring
    ...

Prevention

When it happens

Trigger: An element from an input bundle is fed to an evaluator type that only supports timers (e.g. a timer-only transform's evaluator), i.e. wiring input PCollections into a transform whose evaluator class doesn't override process_element.

Common situations: Building a custom transform/evaluator and forgetting to override process_element; internal wiring mistakes where a data input reaches a timer-handling evaluator; framework-level misuse during custom runner extension.

Related errors


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

Appendix: source

Thrown at sdks/python/apache_beam/runners/direct/transform_evaluator.py:317

    """
    state = self._step_context.get_keyed_state(timer_firing.encoded_key)
    state.clear_timer(
        timer_firing.window,
        timer_firing.name,
        timer_firing.time_domain,
        dynamic_timer_tag=timer_firing.dynamic_timer_tag)
    self.process_timer(timer_firing)

  def process_timer(self, timer_firing):
    """Default process_timer() impl. generating KeyedWorkItem element."""
    self.process_element(
        GlobalWindows.windowed_value(
            KeyedWorkItem(
                timer_firing.encoded_key, timer_firings=[timer_firing])))

  def process_element(self, element):
    """Processes a new element as part of the current bundle."""
    raise NotImplementedError('%s do not process elements.' % type(self))

  def finish_bundle(self) -> TransformResult:
    """Finishes the bundle and produces output."""
    pass


class _BoundedReadEvaluator(_TransformEvaluator):
  """TransformEvaluator for bounded Read transform."""

  # After some benchmarks, 1000 was optimal among {100,1000,10000}
  MAX_ELEMENT_PER_BUNDLE = 1000

  def __init__(
      self,
      evaluation_context,
      applied_ptransform,
      input_committed_bundle,
      side_inputs):

View on GitHub (pinned to 12126d8942)