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
- Override process_element in your evaluator subclass if it should consume elements.
- Don't connect an input PCollection to a transform that only processes timers.
- 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
- Override process_element in any custom evaluator meant to consume elements.
- Don't connect data inputs to timer-only transforms.
- Unit-test evaluators with both element and timer inputs.
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
- DirectRunner: id_label is not supported for PubSub reads
- Execution of [ ] not implemented in runner .
- Root provider for [ ] not implemented in runner
- A BigQuery table or a query must be specified
- A cluster_identifier should be Optional[Union[str…
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)