Join the SIG TFX-Addons community and help make TFX even better!


Wrap node as beam DoFn.

pipeline_node The specification of the node that this launcher lauches.
mlmd_connection_config ML metadata connection config.
pipeline_info The information of the pipeline that this node runs in.
pipeline_runtime_spec The runtime information of the pipeline that this node runs in.
executor_spec Specification for the executor of the node. This is expected for all nodes. This will be used to determine the specific ExecutorOperator class to be used to execute and will be passed into ExecutorOperator.
custom_driver_spec Specification for custom driver. This is expected only for advanced use cases.
deployment_config Deployment Config for the pipeline.

Child Classes

class BundleFinalizerParam

class RestrictionParam

class StateParam

class TimerParam

class WatermarkEstimatorParam





Returns the display data associated to a pipeline component.

It should be reimplemented in pipeline components that wish to have static display data.

Dict[str, Any]: A dictionary containing key:value pairs. The value might be an integer, float or string value; a :class:DisplayDataItem for values that have more data (e.g. short value, label, url); or a :class:HasDisplayData instance that has more display data that should be picked up. For example::

{ 'key1': 'string_value', 'key2': 1234, 'key3': 3.14159265, 'key4': DisplayDataItem('', url=''), 'key5': subComponent }


Called after a bundle of elements is processed on a worker.



Converts from an FunctionSpec to a Fn object.

Prefer registering a urn with its parameter type and constructor.



Gets and/or initializes type hints for this object.

If type hints have not been set, attempts to initialize type hints in this order:

  • Using self.default_type_hints().
  • Using self.class type hints.



View source

Executes node based on signals.

element a signal element to trigger the node.
*signals side input signals indicate completeness of upstream nodes.


Registers and implements the given urn via pickling.


Registers a urn with a constructor.

For example, if 'beam:fn:foo' had parameter type FooPayload, one could write RunnerApiFn.register_urn('bean:fn:foo', FooPayload, foo_from_proto) where foo_from_proto took as arguments a FooPayload and a PipelineContext. This function can also be used as a decorator rather than passing the callable in as the final parameter.

A corresponding to_runner_api_parameter method would be expected that returns the tuple ('beam:fn:foo', FooPayload)


Called to prepare an instance for processing bundles of elements.

This is a good place to initialize transient in-memory resources, such as network connections. The resources can then be disposed in DoFn.teardown.


Called before a bundle of elements is processed on a worker.

Elements to be processed are split into bundles and distributed to workers. Before a worker calls process() on the first element of its bundle, it calls this method.


Called to use to clean up this instance before it is discarded.

A runner will do its best to call this method on any given instance to prevent leaks of transient resources, however, there may be situations where this is impossible (e.g. process crash, hardware failure, etc.) or unnecessary (e.g. the pipeline is shutting down and the process is about to be killed anyway, so all transient resources will be released automatically by the OS). In these cases, the call may not happen. It will also not be retried, because in such situations the DoFn instance no longer exists, so there's no instance to retry it on.

Thus, all work that depends on input elements, and all externally important side effects, must be performed in DoFn.process or DoFn.finish_bundle.


Returns an FunctionSpec encoding this Fn.

Prefer overriding self.to_runner_api_parameter.



A decorator on process fn specifying that the fn performs an unbounded amount of work per input element.



DoFnProcessParams [ElementParam, SideInputParam, TimestampParam, WindowParam, <class 'apache_beam.transforms.core._WatermarkEstimatorParam'>, PaneInfoParam, <class 'apache_beam.transforms.core._BundleFinalizerParam'>, KeyParam, <class 'apache_beam.transforms.core._StateDoFnParam'>, <class 'apache_beam.transforms.core._TimerDoFnParam'>]
DynamicTimerTagParam Instance of apache_beam.transforms.core._DoFnParam
ElementParam Instance of apache_beam.transforms.core._DoFnParam
KeyParam Instance of apache_beam.transforms.core._DoFnParam
PaneInfoParam Instance of apache_beam.transforms.core._DoFnParam
SideInputParam Instance of apache_beam.transforms.core._DoFnParam
TimestampParam Instance of apache_beam.transforms.core._DoFnParam
WindowParam Instance of apache_beam.transforms.core._DoFnParam