Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions sdks/python/apache_beam/yaml/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,17 @@ def _preparse_jinja_flags(argv):
return argv

jinja_variable_parser = argparse.ArgumentParser(allow_abbrev=False)
# Guard against jinja_variable_flags colliding with pipeline options.
# If a flag collides with a known pipeline option, skip it and require
# the variable to be provided via --jinja_variables JSON instead.
try:
from apache_beam.options.pipeline_options import PipelineOptions
_pipeline_option_names = set(PipelineOptions([]).get_all_options().keys())
except Exception:
_pipeline_option_names = set()
for flag_name in jinja_args.jinja_variable_flags:
if flag_name.replace('-', '_') in _pipeline_option_names:
continue
jinja_variable_parser.add_argument('--' + flag_name)
jinja_flag_variables, pipeline_args = jinja_variable_parser.parse_known_args(
other_args)
Expand Down
14 changes: 14 additions & 0 deletions sdks/python/apache_beam/yaml/main_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,20 @@ def test_preparse_jinja_flags(self):
'pos_arg',
])

def test_preparse_jinja_flags_pipeline_option_collision(self):
# A jinja_variable_flags entry that collides with a known pipeline
# option (e.g. runner) must not swallow the pipeline flag.
argv = [
'--jinja_variable_flags=runner,var',
'--runner=DirectRunner',
'--var=my_line',
]
self.assertCountEqual(
main._preparse_jinja_flags(argv), [
'--runner=DirectRunner',
'--jinja_variables=' + '{"var": "my_line"}',
])

def test_jinja_datetime(self):
with tempfile.TemporaryDirectory() as tmpdir:
out_path = os.path.join(tmpdir, 'out.txt')
Expand Down
Loading