Skip to content

Commit e64815d

Browse files
Hint Suggestions for invalid pipeline options (#36072)
* Hint Suggestions for invalid pipeline options * only show suggestions once * Update sdks/python/apache_beam/options/pipeline_options.py Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --------- Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
1 parent becfcf8 commit e64815d

2 files changed

Lines changed: 29 additions & 6 deletions

File tree

sdks/python/apache_beam/options/pipeline_options.py

Lines changed: 25 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
# pytype: skip-file
2121

2222
import argparse
23+
import difflib
2324
import json
2425
import logging
2526
import os
@@ -449,11 +450,30 @@ def from_dictionary(cls, options):
449450

450451
return cls(flags)
451452

453+
@staticmethod
454+
def _warn_on_unknown_options(unknown_args, parser):
455+
if not unknown_args:
456+
return
457+
458+
all_known_options = [
459+
opt for action in parser._actions for opt in action.option_strings
460+
]
461+
462+
for arg in unknown_args:
463+
msg = f"Unparseable argument: {arg}"
464+
if arg.startswith('--'):
465+
arg_name = arg.split('=', 1)[0]
466+
suggestions = difflib.get_close_matches(arg_name, all_known_options)
467+
if suggestions:
468+
msg += f". Did you mean '{suggestions[0]}'?'"
469+
_LOGGER.warning(msg)
470+
452471
def get_all_options(
453472
self,
454473
drop_default=False,
455474
add_extra_args_fn: Optional[Callable[[_BeamArgumentParser], None]] = None,
456-
retain_unknown_options=False) -> Dict[str, Any]:
475+
retain_unknown_options=False,
476+
display_warnings=False) -> Dict[str, Any]:
457477
"""Returns a dictionary of all defined arguments.
458478
459479
Returns a dictionary of all defined arguments (arguments that are defined in
@@ -485,12 +505,11 @@ def get_all_options(
485505
add_extra_args_fn(parser)
486506

487507
known_args, unknown_args = parser.parse_known_args(self._flags)
488-
if retain_unknown_options:
489-
if unknown_args:
490-
_LOGGER.warning(
491-
'Unknown pipeline options received: %s. Ignore if flags are '
492-
'used for internal purposes.' % (','.join(unknown_args)))
493508

509+
if display_warnings:
510+
self._warn_on_unknown_options(unknown_args, parser)
511+
512+
if retain_unknown_options:
494513
seen = set()
495514

496515
def add_new_arg(arg, **kwargs):

sdks/python/apache_beam/pipeline.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -580,6 +580,10 @@ def run(self, test_runner_api='AUTO'):
580580
# type: (Union[bool, str]) -> PipelineResult
581581

582582
"""Runs the pipeline. Returns whatever our runner returns after running."""
583+
# All pipeline options are finalized at this point.
584+
# Call get_all_options to print warnings on invalid options.
585+
self.options.get_all_options(
586+
retain_unknown_options=True, display_warnings=True)
583587

584588
for error_handler in self._error_handlers:
585589
error_handler.verify_closed()

0 commit comments

Comments
 (0)