This is an automated email from the ASF dual-hosted git repository.
jrmccluskey pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 55fde074af8 Add the directory with staged files to sys.path and
document the usage (#39434)
55fde074af8 is described below
commit 55fde074af8891ed6ec1c0bd21454009d6ccc0d5
Author: tvalentyn <[email protected]>
AuthorDate: Tue Jul 28 06:55:00 2026 -0700
Add the directory with staged files to sys.path and document the usage
(#39434)
of `--files_to_stage` and `--beam_plugins`
---
.github/trigger_files/beam_PostCommit_Python.json | 2 +-
CHANGES.md | 2 +-
.../python/apache_beam/options/pipeline_options.py | 1 +
.../apache_beam/options/pipeline_options_test.py | 16 +++++
.../runners/portability/beam_plugins_it_test.py | 70 ++++++++++++++++++++++
.../apache_beam/runners/worker/sdk_worker_main.py | 7 +++
.../runners/worker/sdk_worker_main_test.py | 49 +++++++++++++++
.../sdks/python-pipeline-dependencies.md | 41 ++++++++++++-
8 files changed, 185 insertions(+), 3 deletions(-)
diff --git a/.github/trigger_files/beam_PostCommit_Python.json
b/.github/trigger_files/beam_PostCommit_Python.json
index 00c0cdc3f9c..89cec619b02 100644
--- a/.github/trigger_files/beam_PostCommit_Python.json
+++ b/.github/trigger_files/beam_PostCommit_Python.json
@@ -1,5 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to
run.",
"pr": "38701",
- "modification": 55
+ "modification": 56
}
diff --git a/CHANGES.md b/CHANGES.md
index 7ea2cccde29..1bbde577e87 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -70,8 +70,8 @@
* (Python) Removed the `envoy-data-plane` (and transitive `betterproto`)
dependency; `EnvoyRateLimiter` now uses a small vendored protobuf definition
instead, resolving dependency conflicts for downstream projects
([#37854](https://github.com/apache/beam/issues/37854)).
* (Java) Supported acknowledge mode for JmsIO
([#39253](https://github.com/apache/beam/issues/39253)).
+* (Python) Staged files directory is now automatically added to `sys.path` on
the Python SDK worker at startup. This makes Python files provided via the
'--files_to_stage' pipeline option importable in the pipeline code and makes it
easier to initialize Python SDK harness at startup via the `--beam_plugins`
pipeline option. For more information, see the [Staging Individual
Files](https://beam.apache.org/documentation/sdks/python-pipeline-dependencies/#staging-files)
section of the depend [...]
* (Python) Added `equal_to_approx`, an `assert_that` matcher that compares
numeric pipeline outputs with a configurable tolerance
([#18028](https://github.com/apache/beam/issues/18028)).
-* X feature added (Java/Python)
([#X](https://github.com/apache/beam/issues/X)).
## Breaking Changes
diff --git a/sdks/python/apache_beam/options/pipeline_options.py
b/sdks/python/apache_beam/options/pipeline_options.py
index afbec3f46f7..239d577cfa0 100644
--- a/sdks/python/apache_beam/options/pipeline_options.py
+++ b/sdks/python/apache_beam/options/pipeline_options.py
@@ -1863,6 +1863,7 @@ class SetupOptions(PipelineOptions):
'workers will install them in same order they were specified on '
'the command line.'))
parser.add_argument(
+ '--file_to_stage',
'--files_to_stage',
dest='files_to_stage',
action='append',
diff --git a/sdks/python/apache_beam/options/pipeline_options_test.py
b/sdks/python/apache_beam/options/pipeline_options_test.py
index 90ad27d0a39..b321314b4b7 100644
--- a/sdks/python/apache_beam/options/pipeline_options_test.py
+++ b/sdks/python/apache_beam/options/pipeline_options_test.py
@@ -569,6 +569,22 @@ class PipelineOptionsTest(unittest.TestCase):
options = PipelineOptions(flags=[''])
self.assertEqual(options.get_all_options()['extra_packages'], None)
+ def test_files_to_stage(self):
+ options = PipelineOptions([
+ '--file_to_stage',
+ 'abc',
+ '--files_to_stage',
+ 'def',
+ '--files_to_stage',
+ 'ghi'
+ ])
+ self.assertEqual(
+ sorted(options.get_all_options()['files_to_stage']),
+ ['abc', 'def', 'ghi'])
+
+ options = PipelineOptions(flags=[''])
+ self.assertEqual(options.get_all_options()['files_to_stage'], None)
+
def test_dataflow_job_file(self):
options = PipelineOptions(['--dataflow_job_file', 'abc'])
self.assertEqual(options.get_all_options()['dataflow_job_file'], 'abc')
diff --git
a/sdks/python/apache_beam/runners/portability/beam_plugins_it_test.py
b/sdks/python/apache_beam/runners/portability/beam_plugins_it_test.py
new file mode 100644
index 00000000000..54984503fa5
--- /dev/null
+++ b/sdks/python/apache_beam/runners/portability/beam_plugins_it_test.py
@@ -0,0 +1,70 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+import logging
+import os
+import shutil
+import tempfile
+import unittest
+import uuid
+
+import pytest
+
+import apache_beam as beam
+from apache_beam.options.pipeline_options import SetupOptions
+from apache_beam.testing.test_pipeline import TestPipeline
+from apache_beam.testing.util import assert_that
+from apache_beam.testing.util import equal_to
+
+_LOGGER = logging.getLogger(__name__)
+
+
+class BeamPluginsIT(unittest.TestCase):
+ def setUp(self):
+ self.temp_dir = tempfile.mkdtemp()
+
+ def tearDown(self):
+ shutil.rmtree(self.temp_dir)
+
+ @pytest.mark.it_postcommit
+ def test_beam_plugins_staging(self):
+ pipeline = TestPipeline(is_integration_test=True)
+ setup_options = pipeline.options.view_as(SetupOptions)
+
+ plugin_name = f'beam_integration_plugin_{uuid.uuid4().hex[:8]}'
+ plugin_file_path = os.path.join(self.temp_dir, f'{plugin_name}.py')
+
+ with open(plugin_file_path, 'w') as f:
+ f.write("import sys\nsys.beam_plugin_loaded_for_test = True\n")
+
+ staged_files = setup_options.files_to_stage or []
+ staged_files.append(plugin_file_path)
+ setup_options.files_to_stage = staged_files
+ setup_options.beam_plugins = [plugin_name]
+
+ def check_plugin_loaded(_):
+ import sys
+ return getattr(sys, 'beam_plugin_loaded_for_test', False)
+
+ with pipeline as p:
+ res = (p | beam.Create([None]) | beam.Map(check_plugin_loaded))
+ assert_that(res, equal_to([True]))
+
+
+if __name__ == '__main__':
+ logging.getLogger().setLevel(logging.INFO)
+ unittest.main()
diff --git a/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
b/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
index 754a631eaf3..1c464141f7e 100644
--- a/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
+++ b/sdks/python/apache_beam/runners/worker/sdk_worker_main.py
@@ -50,6 +50,7 @@ from apache_beam.utils import profiler
_LOGGER = logging.getLogger(__name__)
_ENABLE_GOOGLE_CLOUD_PROFILER = 'enable_google_cloud_profiler'
_FN_LOG_HANDLER = None
+_STAGED_DIRECTORY = 'staged'
def _import_beam_plugins(plugins):
@@ -131,6 +132,12 @@ def create_harness(environment, dry_run=False):
environment.get('RUNNER_CAPABILITIES', '').split())
_LOGGER.info('semi_persistent_directory: %s', semi_persistent_directory)
+ experiments = sdk_pipeline_options.view_as(DebugOptions).experiments or []
+ if 'no_staged_dir_in_sys_path' not in experiments and
semi_persistent_directory:
+ staged_dir = os.path.join(semi_persistent_directory, _STAGED_DIRECTORY)
+ if os.path.isdir(staged_dir) and staged_dir not in sys.path:
+ sys.path.append(staged_dir)
+
_worker_id = environment.get('WORKER_ID', None)
try:
diff --git a/sdks/python/apache_beam/runners/worker/sdk_worker_main_test.py
b/sdks/python/apache_beam/runners/worker/sdk_worker_main_test.py
index 5ecd9616fcf..c84438953b1 100644
--- a/sdks/python/apache_beam/runners/worker/sdk_worker_main_test.py
+++ b/sdks/python/apache_beam/runners/worker/sdk_worker_main_test.py
@@ -129,6 +129,55 @@ class SdkWorkerMainTest(unittest.TestCase):
def test_import_beam_plugins(self):
sdk_worker_main._import_beam_plugins(BeamPlugin.get_all_plugin_paths())
+ def test_create_harness_adds_staged_dir_to_sys_path(self):
+ import sys
+ import tempfile
+
+ with tempfile.TemporaryDirectory() as temp_dir:
+ staged_dir = os.path.join(temp_dir, sdk_worker_main._STAGED_DIRECTORY)
+ os.mkdir(staged_dir)
+
+ env = {
+ 'CONTROL_API_SERVICE_DESCRIPTOR': '',
+ 'SEMI_PERSISTENT_DIRECTORY': temp_dir,
+ }
+
+ if staged_dir in sys.path:
+ sys.path.remove(staged_dir)
+
+ sdk_worker_main.create_harness(env, dry_run=True)
+
+ try:
+ self.assertIn(staged_dir, sys.path)
+ finally:
+ if staged_dir in sys.path:
+ sys.path.remove(staged_dir)
+
+ def test_create_harness_does_not_add_staged_dir_with_experiment(self):
+ import sys
+ import tempfile
+
+ with tempfile.TemporaryDirectory() as temp_dir:
+ staged_dir = os.path.join(temp_dir, sdk_worker_main._STAGED_DIRECTORY)
+ os.mkdir(staged_dir)
+
+ env = {
+ 'CONTROL_API_SERVICE_DESCRIPTOR': '',
+ 'SEMI_PERSISTENT_DIRECTORY': temp_dir,
+ 'PIPELINE_OPTIONS': '{"experiments":["no_staged_dir_in_sys_path"]}',
+ }
+
+ if staged_dir in sys.path:
+ sys.path.remove(staged_dir)
+
+ sdk_worker_main.create_harness(env, dry_run=True)
+
+ try:
+ self.assertNotIn(staged_dir, sys.path)
+ finally:
+ if staged_dir in sys.path:
+ sys.path.remove(staged_dir)
+
@staticmethod
def _overrides_case_to_option_dict(case):
"""
diff --git
a/website/www/site/content/en/documentation/sdks/python-pipeline-dependencies.md
b/website/www/site/content/en/documentation/sdks/python-pipeline-dependencies.md
index 26dc1db1bc0..7d27f91d322 100644
---
a/website/www/site/content/en/documentation/sdks/python-pipeline-dependencies.md
+++
b/website/www/site/content/en/documentation/sdks/python-pipeline-dependencies.md
@@ -91,9 +91,48 @@ If your pipeline uses packages that are not available
publicly (e.g. packages th
See the [build
documentation](https://pypa-build.readthedocs.io/en/latest/index.html) for more
details on this command.
+## Staging Individual Files {#staging-files}
+
+If your pipeline relies on one or more individual Python files or non-Python
data files that do not need to be packaged as a full Python package, you can
stage them individually to the remote workers.
+
+To stage individual files, run your pipeline with the `--files_to_stage` (or
`--file_to_stage`) pipeline option. This option accepts a list of local file
paths:
+
+ --files_to_stage="/path/to/my_module.py,/path/to/data_config.json"
+
+When the pipeline runs, the runner uploads these files and makes them
available on the workers in the worker's staged files directory.
+
+### Accessing Staged Files on the Workers
+
+Staged files are downloaded onto the worker:
+
+* **Python modules**: Starting with Apache Beam 2.76.0, to make it easy to
import staged Python files as modules, Beam automatically appends the worker's
staged files directory to the Python interpreter's search path (`sys.path`)
during startup. This means you can import them directly in your code:
+
+ import my_module
+
+* **Data or configuration files**: If you staged non-Python files (such as a
JSON config), they are downloaded to the staged files directory. You can locate
this directory on the worker by reading the `SEMI_PERSISTENT_DIRECTORY`
environment variable, or look for `/tmp/staged` which is the default location
for staged files on containerized runners:
+
+ import os
+ staged_dir = os.environ.get('SEMI_PERSISTENT_DIRECTORY', '/tmp/staged')
+ config_path = os.path.join(staged_dir, 'data_config.json')
+
+### Importing Plugins on Worker Startup
+
+You can use staged files in combination with the `--beam_plugins` pipeline
option (supported starting with Apache Beam 2.76.0) to run initialization code
on the workers before any processing starts.
+
+To use a staged file as a plugin:
+1. Stage the plugin file (e.g. `my_custom_plugin.py`) using `--file_to_stage`.
+2. Reference the module name in `--beam_plugins`.
+
+For example, run your pipeline with:
+
+ --file_to_stage="/path/to/my_custom_plugin.py" \
+ --beam_plugins="my_custom_plugin"
+
+This instructs the SDK worker to import `my_custom_plugin` immediately on
startup, triggering any initialization logic defined in the module.
+
## Multiple File Dependencies {#multiple-file-dependencies}
-Often, your pipeline code spans multiple files. To run your project remotely,
you must group these files as a Python package and specify the package when you
run your pipeline. When the remote workers start, they will install your
package. To group your files as a Python package and make it available
remotely, perform the following steps:
+Often, your pipeline code spans multiple files. To run your project remotely,
it is best to group these files as a Python package and specify the package
when you run your pipeline. When the remote workers start, they will install
your package. To group your files as a Python package and make it available
remotely, perform the following steps:
1. Create a
[setup.py](https://pythonhosted.org/an_example_pypi_project/setuptools.html)
file for your project. The following is a very basic `setup.py` file.