apache_beam.runners package
Runner objects execute a Pipeline.
This package defines runners, which are used to execute a pipeline.
Subpackages
- apache_beam.runners.dask package
- apache_beam.runners.dataflow package
- Submodules
- apache_beam.runners.dataflow.dataflow_exercise_metrics_pipeline module
- apache_beam.runners.dataflow.dataflow_exercise_streaming_metrics_pipeline module
- apache_beam.runners.dataflow.dataflow_job_service module
- apache_beam.runners.dataflow.dataflow_metrics module
- apache_beam.runners.dataflow.dataflow_runner module
- apache_beam.runners.dataflow.ptransform_overrides module
- apache_beam.runners.dataflow.test_dataflow_runner module
- Submodules
- apache_beam.runners.direct package
- Submodules
- apache_beam.runners.direct.bundle_factory module
- apache_beam.runners.direct.clock module
- apache_beam.runners.direct.consumer_tracking_pipeline_visitor module
- apache_beam.runners.direct.direct_metrics module
- apache_beam.runners.direct.direct_runner module
- apache_beam.runners.direct.direct_userstate module
- apache_beam.runners.direct.evaluation_context module
- apache_beam.runners.direct.executor module
- apache_beam.runners.direct.sdf_direct_runner module
- apache_beam.runners.direct.test_direct_runner module
- apache_beam.runners.direct.test_stream_impl module
- apache_beam.runners.direct.transform_evaluator module
- apache_beam.runners.direct.util module
- apache_beam.runners.direct.watermark_manager module
- Submodules
- apache_beam.runners.interactive package
- Subpackages
- apache_beam.runners.interactive.caching package
- apache_beam.runners.interactive.dataproc package
- apache_beam.runners.interactive.display package
- apache_beam.runners.interactive.messaging package
- apache_beam.runners.interactive.options package
- apache_beam.runners.interactive.sql package
- apache_beam.runners.interactive.testing package
- Submodules
- apache_beam.runners.interactive.augmented_pipeline module
- apache_beam.runners.interactive.background_caching_job module
BackgroundCachingJobattempt_to_run_background_caching_job()is_background_caching_job_needed()is_cache_complete()has_source_to_cache()attempt_to_cancel_background_caching_job()attempt_to_stop_test_stream_service()is_a_test_stream_service_running()is_source_to_cache_changed()extract_source_to_cache_signature()
- apache_beam.runners.interactive.cache_manager module
- apache_beam.runners.interactive.interactive_beam module
- apache_beam.runners.interactive.interactive_environment module
- apache_beam.runners.interactive.interactive_runner module
- apache_beam.runners.interactive.pipeline_fragment module
- apache_beam.runners.interactive.pipeline_instrument module
- apache_beam.runners.interactive.recording_manager module
- apache_beam.runners.interactive.user_pipeline_tracker module
- apache_beam.runners.interactive.utils module
to_element_list()elements_to_df()register_ipython_log_handler()IPythonLogHandlerobfuscate()ProgressIndicatorprogress_indicated()as_json()deferred_df_to_pcollection()pcoll_by_name()find_pcoll_name()cacheables()watch_sources()has_unbounded_sources()unbounded_sources()create_var_in_main()assert_bucket_exists()detect_pipeline_runner()
- Subpackages
- apache_beam.runners.job package
Submodules
- apache_beam.runners.pipeline_context module
PortableObjectPipelineContextPipelineContext.add_requirement()PipelineContext.requirements()PipelineContext.coder_id_from_element_type()PipelineContext.deterministic_coder()PipelineContext.element_type_from_coder_id()PipelineContext.from_runner_api()PipelineContext.to_runner_api()PipelineContext.default_environment_id()PipelineContext.get_environment_id_for_resource_hints()
- apache_beam.runners.pipeline_utils module
- apache_beam.runners.render module
RenderOptionsPipelineRendererPipelineRenderer.update()PipelineRenderer.style()PipelineRenderer.to_dot()PipelineRenderer.to_dot_iter()PipelineRenderer.transform_to_dot()PipelineRenderer.transform_node()PipelineRenderer.transform_attributes()PipelineRenderer.pcoll_leaf_consumers_iter()PipelineRenderer.pcoll_leaf_consumers()PipelineRenderer.is_leaf()PipelineRenderer.info()PipelineRenderer.layout_dot()PipelineRenderer.page_callback_data()PipelineRenderer.render_data()PipelineRenderer.render_json()PipelineRenderer.page()
RenderRunnerRenderPipelineResultrun()render_one()run_server()
- apache_beam.runners.runner module
PipelineRunnerPipelineRunner.run()PipelineRunner.run_async()PipelineRunner.run_portable_pipeline()PipelineRunner.default_environment()PipelineRunner.run_pipeline()PipelineRunner.apply()PipelineRunner.apply_PTransform()PipelineRunner.is_fnapi_compatible()PipelineRunner.check_requirements()PipelineRunner.default_pickle_library_override()
PipelineStatePipelineState.UNKNOWNPipelineState.STARTINGPipelineState.STOPPEDPipelineState.RUNNINGPipelineState.DONEPipelineState.FAILEDPipelineState.CANCELLEDPipelineState.UPDATEDPipelineState.DRAININGPipelineState.DRAINEDPipelineState.PENDINGPipelineState.CANCELLINGPipelineState.RESOURCE_CLEANING_UPPipelineState.UNRECOGNIZEDPipelineState.is_terminal()
PipelineResult
- apache_beam.runners.sdf_utils module
SplitResultPrimarySplitResultResidualThreadsafeRestrictionTrackerThreadsafeRestrictionTracker.current_restriction()ThreadsafeRestrictionTracker.try_claim()ThreadsafeRestrictionTracker.defer_remainder()ThreadsafeRestrictionTracker.check_done()ThreadsafeRestrictionTracker.current_progress()ThreadsafeRestrictionTracker.try_split()ThreadsafeRestrictionTracker.deferred_status()ThreadsafeRestrictionTracker.is_bounded()
RestrictionTrackerViewThreadsafeWatermarkEstimatorNoOpWatermarkEstimatorProvider
- apache_beam.runners.trivial_runner module