From 7b16a5b0042908318193f40f58c562cb24a32d25 Mon Sep 17 00:00:00 2001 From: Shubham Kumar Date: Thu, 21 May 2026 10:51:37 +0000 Subject: [PATCH 1/3] adding a flag to log and skip a file if grib-filter-expression is invalid --- .../splitter_pipeline/file_splitters.py | 104 ++++++++++-------- weather_sp/splitter_pipeline/pipeline.py | 12 +- 2 files changed, 69 insertions(+), 47 deletions(-) diff --git a/weather_sp/splitter_pipeline/file_splitters.py b/weather_sp/splitter_pipeline/file_splitters.py index f3d18b4..a128524 100644 --- a/weather_sp/splitter_pipeline/file_splitters.py +++ b/weather_sp/splitter_pipeline/file_splitters.py @@ -72,7 +72,7 @@ class FileSplitter(abc.ABC): def __init__(self, input_path: str, output_info: OutFileInfo, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None): + grib_filter_expression: t.Optional[str] = None, skip_on_invalid_grib_filter_expression:bool = False): self.input_path = input_path self.output_info = output_info self.force_split = force_split @@ -81,6 +81,7 @@ def __init__(self, input_path: str, output_info: OutFileInfo, self.logger.debug('Splitter for path=%s, output base=%s', self.input_path, self.output_info) self.grib_filter_expression = grib_filter_expression + self.skip_on_invalid_grib_filter_expression = skip_on_invalid_grib_filter_expression @abc.abstractmethod def split_data(self) -> None: @@ -207,7 +208,7 @@ def split_data(self) -> None: grib_copy_cmd = shutil.which('grib_copy') grib_get_cmd = shutil.which('grib_get') uniq_cmd = shutil.which('uniq') - for cmd, name in [(grib_get_cmd, 'grib_copy'), (grib_get_cmd, 'grib_get'), (uniq_cmd, 'uniq')]: + for cmd, name in [(grib_copy_cmd, 'grib_copy'), (grib_get_cmd, 'grib_get'), (uniq_cmd, 'uniq')]: if not cmd: raise EnvironmentError(f'binary {name!r} is not available in the current environment!') @@ -227,49 +228,61 @@ def split_data(self) -> None: # This ensures dims like time are represented as 0600 instead of 600. split_dims_arg = ','.join(f'{dim}:s' for dim in split_dims) with self._copy_to_local_file() as local_file: - self.logger.info('Skipping as needed...') - # Append -w flag to filter GRIB messages matching the given expression - if self.grib_filter_expression: - grib_get_args = [grib_get_cmd, '-p', split_dims_arg, '-w', self.grib_filter_expression, local_file.name] - else: - grib_get_args = [grib_get_cmd, '-p', split_dims_arg, local_file.name] - grib_get_process = subprocess.Popen(grib_get_args, stdout=subprocess.PIPE) - uniq_output = subprocess.check_output((uniq_cmd,), stdin=grib_get_process.stdout) - output_paths = [] - skipped_paths = [] - for line in uniq_output.decode('utf-8').rstrip('\n').split('\n'): - splits = dict(zip(split_dims, line.split(' '))) - output_path = self.output_info.formatted_output_path(splits) - if self.should_skip_file(output_path): - skipped_paths.append(output_path) - continue - output_paths.append(output_path) - if not output_paths: - metrics.Metrics.counter('file_splitters', 'skipped').inc() - self.logger.info('Skipping %s, file already split into: %s', - repr(self.input_path), ', '.join(skipped_paths)) - return - - with tempfile.TemporaryDirectory() as tmpdir: - self.logger.info('Performing split.') - dest = os.path.join(tmpdir, flat_output_template) + try: + self.logger.info('Skipping as needed...') + # Append -w flag to filter GRIB messages matching the given expression if self.grib_filter_expression: - subprocess.run([grib_copy_cmd, "-w", - self.grib_filter_expression, - local_file.name, dest], check=True) + grib_get_args = [grib_get_cmd, '-p', split_dims_arg, '-w', self.grib_filter_expression, local_file.name] else: - subprocess.run([grib_copy_cmd, local_file.name, dest], - check=True) - - self.logger.info('Uploading %r...', self.input_path) - for flat_target in os.listdir(tmpdir): - dest_file_path = f'{prefix}{flat_target.replace(delimiter, slash)}' - self.logger.info([prefix, dest_file_path, local_file.name, - self.output_info.unformatted_output_path()]) - - copy(os.path.join(tmpdir, flat_target), dest_file_path) - self.logger.info('Finished uploading %r', self.input_path) - + grib_get_args = [grib_get_cmd, '-p', split_dims_arg, local_file.name] + grib_get_process = subprocess.Popen(grib_get_args, stdout=subprocess.PIPE) + uniq_output = subprocess.check_output((uniq_cmd,), stdin=grib_get_process.stdout) + output_paths = [] + skipped_paths = [] + for line in uniq_output.decode('utf-8').rstrip('\n').split('\n'): + splits = dict(zip(split_dims, line.split(' '))) + output_path = self.output_info.formatted_output_path(splits) + if self.should_skip_file(output_path): + skipped_paths.append(output_path) + continue + output_paths.append(output_path) + if not output_paths: + metrics.Metrics.counter('file_splitters', 'skipped').inc() + self.logger.info('Skipping %s, file already split into: %s', + repr(self.input_path), ', '.join(skipped_paths)) + return + + with tempfile.TemporaryDirectory() as tmpdir: + self.logger.info('Performing split.') + dest = os.path.join(tmpdir, flat_output_template) + if self.grib_filter_expression: + subprocess.run([grib_copy_cmd, "-w", + self.grib_filter_expression, + local_file.name, dest], check=True) + else: + subprocess.run([grib_copy_cmd, local_file.name, dest], + check=True) + + self.logger.info('Uploading %r...', self.input_path) + for flat_target in os.listdir(tmpdir): + dest_file_path = f'{prefix}{flat_target.replace(delimiter, slash)}' + self.logger.info([prefix, dest_file_path, local_file.name, + self.output_info.unformatted_output_path()]) + + copy(os.path.join(tmpdir, flat_target), dest_file_path) + self.logger.info('Finished uploading %r', self.input_path) + except Exception as e: + log_msg = ( + f"GRIB tool failed for {self.input_path!r}. This means the requested " + f"filter expression {self.grib_filter_expression!r} does not exist in this file. " + f"Error: {e}" + ) + if self.skip_on_invalid_grib_filter_expression: + self.logger.warning(f"{log_msg} | Flag 'skip_on_invalid_grib_filter_expression' is True. Skipping file.") + return + else: + self.logger.error(f"{log_msg} | Flag 'skip_on_invalid_grib_filter_expression' is False. Error raised") + raise class NetCdfSplitter(FileSplitter): @@ -353,7 +366,8 @@ def get_splitter(file_path: str, dry_run: bool, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None) -> FileSplitter: + grib_filter_expression: t.Optional[str] = None, + skip_on_invalid_grib_filter_expression: bool = False) -> FileSplitter: if dry_run: logger.info('Using splitter: DrySplitter') return DrySplitter(file_path, output_info, logging_level=logging_level) @@ -371,7 +385,7 @@ def get_splitter(file_path: str, if cmd: logger.info('Using splitter: GribSplitterV2') return GribSplitterV2(file_path, output_info, force_split, - logging_level, grib_filter_expression) + logging_level, grib_filter_expression, skip_on_invalid_grib_filter_expression) else: logger.info('Using splitter: GribSplitter') return GribSplitter(file_path, output_info, force_split, diff --git a/weather_sp/splitter_pipeline/pipeline.py b/weather_sp/splitter_pipeline/pipeline.py index 17b3251..7eb4547 100644 --- a/weather_sp/splitter_pipeline/pipeline.py +++ b/weather_sp/splitter_pipeline/pipeline.py @@ -46,7 +46,8 @@ def split_file(input_file: str, dry_run: bool, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None): + grib_filter_expression: t.Optional[str] = None, + skip_on_invalid_grib_filter_expression: bool = False): output_base_name = get_output_base_name(input_path=input_file, input_base=input_base_dir, output_template=output_template, @@ -61,7 +62,8 @@ def split_file(input_file: str, dry_run, force_split, level, - grib_filter_expression) + grib_filter_expression, + skip_on_invalid_grib_filter_expression) splitter.split_data() @@ -134,6 +136,10 @@ def run(argv: t.List[str], save_main_session: bool = True): 'specifically supported by the GribSplitterV2' 'implementation.' 'Example: typeOfLevel=isobaricInhPa,level=1000') + parser.add_argument('--skip-on-invalid-grib-filter-expression', action='store_true', default=False, + help='If provided (True), files that do not contain the key-values specified ' + 'in the --where filter will be logged and skipped. By default (False), ' + 'the pipeline will raise an error and break.') parser.add_argument('--topic', type=str, default=None, help='Pub/Sub topic to read from for streaming mode.') parser.add_argument('--subscription', type=str, default=None, @@ -162,6 +168,7 @@ def run(argv: t.List[str], save_main_session: bool = True): formatting = known_args.formatting dry_run = known_args.dry_run grib_filter_expression = known_args.where + skip_on_invalid_grib_filter_expression = known_args.skip_on_invalid_grib_filter_expression if not output_template and not output_dir: raise ValueError('No output specified') @@ -214,6 +221,7 @@ def run(argv: t.List[str], save_main_session: bool = True): known_args.force, known_args.log_level, grib_filter_expression, + skip_on_invalid_grib_filter_expression, ) ) From 47094f37f279f67b99ff267047fb4578db5c3ff2 Mon Sep 17 00:00:00 2001 From: Shubham Kumar Date: Fri, 22 May 2026 08:55:30 +0000 Subject: [PATCH 2/3] Revert "adding a flag to log and skip a file if grib-filter-expression is invalid" This reverts commit 7b16a5b0042908318193f40f58c562cb24a32d25. --- .../splitter_pipeline/file_splitters.py | 104 ++++++++---------- weather_sp/splitter_pipeline/pipeline.py | 12 +- 2 files changed, 47 insertions(+), 69 deletions(-) diff --git a/weather_sp/splitter_pipeline/file_splitters.py b/weather_sp/splitter_pipeline/file_splitters.py index a128524..f3d18b4 100644 --- a/weather_sp/splitter_pipeline/file_splitters.py +++ b/weather_sp/splitter_pipeline/file_splitters.py @@ -72,7 +72,7 @@ class FileSplitter(abc.ABC): def __init__(self, input_path: str, output_info: OutFileInfo, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None, skip_on_invalid_grib_filter_expression:bool = False): + grib_filter_expression: t.Optional[str] = None): self.input_path = input_path self.output_info = output_info self.force_split = force_split @@ -81,7 +81,6 @@ def __init__(self, input_path: str, output_info: OutFileInfo, self.logger.debug('Splitter for path=%s, output base=%s', self.input_path, self.output_info) self.grib_filter_expression = grib_filter_expression - self.skip_on_invalid_grib_filter_expression = skip_on_invalid_grib_filter_expression @abc.abstractmethod def split_data(self) -> None: @@ -208,7 +207,7 @@ def split_data(self) -> None: grib_copy_cmd = shutil.which('grib_copy') grib_get_cmd = shutil.which('grib_get') uniq_cmd = shutil.which('uniq') - for cmd, name in [(grib_copy_cmd, 'grib_copy'), (grib_get_cmd, 'grib_get'), (uniq_cmd, 'uniq')]: + for cmd, name in [(grib_get_cmd, 'grib_copy'), (grib_get_cmd, 'grib_get'), (uniq_cmd, 'uniq')]: if not cmd: raise EnvironmentError(f'binary {name!r} is not available in the current environment!') @@ -228,61 +227,49 @@ def split_data(self) -> None: # This ensures dims like time are represented as 0600 instead of 600. split_dims_arg = ','.join(f'{dim}:s' for dim in split_dims) with self._copy_to_local_file() as local_file: - try: - self.logger.info('Skipping as needed...') - # Append -w flag to filter GRIB messages matching the given expression + self.logger.info('Skipping as needed...') + # Append -w flag to filter GRIB messages matching the given expression + if self.grib_filter_expression: + grib_get_args = [grib_get_cmd, '-p', split_dims_arg, '-w', self.grib_filter_expression, local_file.name] + else: + grib_get_args = [grib_get_cmd, '-p', split_dims_arg, local_file.name] + grib_get_process = subprocess.Popen(grib_get_args, stdout=subprocess.PIPE) + uniq_output = subprocess.check_output((uniq_cmd,), stdin=grib_get_process.stdout) + output_paths = [] + skipped_paths = [] + for line in uniq_output.decode('utf-8').rstrip('\n').split('\n'): + splits = dict(zip(split_dims, line.split(' '))) + output_path = self.output_info.formatted_output_path(splits) + if self.should_skip_file(output_path): + skipped_paths.append(output_path) + continue + output_paths.append(output_path) + if not output_paths: + metrics.Metrics.counter('file_splitters', 'skipped').inc() + self.logger.info('Skipping %s, file already split into: %s', + repr(self.input_path), ', '.join(skipped_paths)) + return + + with tempfile.TemporaryDirectory() as tmpdir: + self.logger.info('Performing split.') + dest = os.path.join(tmpdir, flat_output_template) if self.grib_filter_expression: - grib_get_args = [grib_get_cmd, '-p', split_dims_arg, '-w', self.grib_filter_expression, local_file.name] + subprocess.run([grib_copy_cmd, "-w", + self.grib_filter_expression, + local_file.name, dest], check=True) else: - grib_get_args = [grib_get_cmd, '-p', split_dims_arg, local_file.name] - grib_get_process = subprocess.Popen(grib_get_args, stdout=subprocess.PIPE) - uniq_output = subprocess.check_output((uniq_cmd,), stdin=grib_get_process.stdout) - output_paths = [] - skipped_paths = [] - for line in uniq_output.decode('utf-8').rstrip('\n').split('\n'): - splits = dict(zip(split_dims, line.split(' '))) - output_path = self.output_info.formatted_output_path(splits) - if self.should_skip_file(output_path): - skipped_paths.append(output_path) - continue - output_paths.append(output_path) - if not output_paths: - metrics.Metrics.counter('file_splitters', 'skipped').inc() - self.logger.info('Skipping %s, file already split into: %s', - repr(self.input_path), ', '.join(skipped_paths)) - return - - with tempfile.TemporaryDirectory() as tmpdir: - self.logger.info('Performing split.') - dest = os.path.join(tmpdir, flat_output_template) - if self.grib_filter_expression: - subprocess.run([grib_copy_cmd, "-w", - self.grib_filter_expression, - local_file.name, dest], check=True) - else: - subprocess.run([grib_copy_cmd, local_file.name, dest], - check=True) - - self.logger.info('Uploading %r...', self.input_path) - for flat_target in os.listdir(tmpdir): - dest_file_path = f'{prefix}{flat_target.replace(delimiter, slash)}' - self.logger.info([prefix, dest_file_path, local_file.name, - self.output_info.unformatted_output_path()]) - - copy(os.path.join(tmpdir, flat_target), dest_file_path) - self.logger.info('Finished uploading %r', self.input_path) - except Exception as e: - log_msg = ( - f"GRIB tool failed for {self.input_path!r}. This means the requested " - f"filter expression {self.grib_filter_expression!r} does not exist in this file. " - f"Error: {e}" - ) - if self.skip_on_invalid_grib_filter_expression: - self.logger.warning(f"{log_msg} | Flag 'skip_on_invalid_grib_filter_expression' is True. Skipping file.") - return - else: - self.logger.error(f"{log_msg} | Flag 'skip_on_invalid_grib_filter_expression' is False. Error raised") - raise + subprocess.run([grib_copy_cmd, local_file.name, dest], + check=True) + + self.logger.info('Uploading %r...', self.input_path) + for flat_target in os.listdir(tmpdir): + dest_file_path = f'{prefix}{flat_target.replace(delimiter, slash)}' + self.logger.info([prefix, dest_file_path, local_file.name, + self.output_info.unformatted_output_path()]) + + copy(os.path.join(tmpdir, flat_target), dest_file_path) + self.logger.info('Finished uploading %r', self.input_path) + class NetCdfSplitter(FileSplitter): @@ -366,8 +353,7 @@ def get_splitter(file_path: str, dry_run: bool, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None, - skip_on_invalid_grib_filter_expression: bool = False) -> FileSplitter: + grib_filter_expression: t.Optional[str] = None) -> FileSplitter: if dry_run: logger.info('Using splitter: DrySplitter') return DrySplitter(file_path, output_info, logging_level=logging_level) @@ -385,7 +371,7 @@ def get_splitter(file_path: str, if cmd: logger.info('Using splitter: GribSplitterV2') return GribSplitterV2(file_path, output_info, force_split, - logging_level, grib_filter_expression, skip_on_invalid_grib_filter_expression) + logging_level, grib_filter_expression) else: logger.info('Using splitter: GribSplitter') return GribSplitter(file_path, output_info, force_split, diff --git a/weather_sp/splitter_pipeline/pipeline.py b/weather_sp/splitter_pipeline/pipeline.py index 7eb4547..17b3251 100644 --- a/weather_sp/splitter_pipeline/pipeline.py +++ b/weather_sp/splitter_pipeline/pipeline.py @@ -46,8 +46,7 @@ def split_file(input_file: str, dry_run: bool, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None, - skip_on_invalid_grib_filter_expression: bool = False): + grib_filter_expression: t.Optional[str] = None): output_base_name = get_output_base_name(input_path=input_file, input_base=input_base_dir, output_template=output_template, @@ -62,8 +61,7 @@ def split_file(input_file: str, dry_run, force_split, level, - grib_filter_expression, - skip_on_invalid_grib_filter_expression) + grib_filter_expression) splitter.split_data() @@ -136,10 +134,6 @@ def run(argv: t.List[str], save_main_session: bool = True): 'specifically supported by the GribSplitterV2' 'implementation.' 'Example: typeOfLevel=isobaricInhPa,level=1000') - parser.add_argument('--skip-on-invalid-grib-filter-expression', action='store_true', default=False, - help='If provided (True), files that do not contain the key-values specified ' - 'in the --where filter will be logged and skipped. By default (False), ' - 'the pipeline will raise an error and break.') parser.add_argument('--topic', type=str, default=None, help='Pub/Sub topic to read from for streaming mode.') parser.add_argument('--subscription', type=str, default=None, @@ -168,7 +162,6 @@ def run(argv: t.List[str], save_main_session: bool = True): formatting = known_args.formatting dry_run = known_args.dry_run grib_filter_expression = known_args.where - skip_on_invalid_grib_filter_expression = known_args.skip_on_invalid_grib_filter_expression if not output_template and not output_dir: raise ValueError('No output specified') @@ -221,7 +214,6 @@ def run(argv: t.List[str], save_main_session: bool = True): known_args.force, known_args.log_level, grib_filter_expression, - skip_on_invalid_grib_filter_expression, ) ) From 81bbaeaa6d4688fe1dd82953f567b263c1b3004a Mon Sep 17 00:00:00 2001 From: Shubham Kumar Date: Fri, 22 May 2026 09:11:56 +0000 Subject: [PATCH 3/3] Reapply "adding a flag to log and skip a file if grib-filter-expression is invalid" This reverts commit 47094f37f279f67b99ff267047fb4578db5c3ff2. --- .../splitter_pipeline/file_splitters.py | 104 ++++++++++-------- weather_sp/splitter_pipeline/pipeline.py | 12 +- 2 files changed, 69 insertions(+), 47 deletions(-) diff --git a/weather_sp/splitter_pipeline/file_splitters.py b/weather_sp/splitter_pipeline/file_splitters.py index f3d18b4..a128524 100644 --- a/weather_sp/splitter_pipeline/file_splitters.py +++ b/weather_sp/splitter_pipeline/file_splitters.py @@ -72,7 +72,7 @@ class FileSplitter(abc.ABC): def __init__(self, input_path: str, output_info: OutFileInfo, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None): + grib_filter_expression: t.Optional[str] = None, skip_on_invalid_grib_filter_expression:bool = False): self.input_path = input_path self.output_info = output_info self.force_split = force_split @@ -81,6 +81,7 @@ def __init__(self, input_path: str, output_info: OutFileInfo, self.logger.debug('Splitter for path=%s, output base=%s', self.input_path, self.output_info) self.grib_filter_expression = grib_filter_expression + self.skip_on_invalid_grib_filter_expression = skip_on_invalid_grib_filter_expression @abc.abstractmethod def split_data(self) -> None: @@ -207,7 +208,7 @@ def split_data(self) -> None: grib_copy_cmd = shutil.which('grib_copy') grib_get_cmd = shutil.which('grib_get') uniq_cmd = shutil.which('uniq') - for cmd, name in [(grib_get_cmd, 'grib_copy'), (grib_get_cmd, 'grib_get'), (uniq_cmd, 'uniq')]: + for cmd, name in [(grib_copy_cmd, 'grib_copy'), (grib_get_cmd, 'grib_get'), (uniq_cmd, 'uniq')]: if not cmd: raise EnvironmentError(f'binary {name!r} is not available in the current environment!') @@ -227,49 +228,61 @@ def split_data(self) -> None: # This ensures dims like time are represented as 0600 instead of 600. split_dims_arg = ','.join(f'{dim}:s' for dim in split_dims) with self._copy_to_local_file() as local_file: - self.logger.info('Skipping as needed...') - # Append -w flag to filter GRIB messages matching the given expression - if self.grib_filter_expression: - grib_get_args = [grib_get_cmd, '-p', split_dims_arg, '-w', self.grib_filter_expression, local_file.name] - else: - grib_get_args = [grib_get_cmd, '-p', split_dims_arg, local_file.name] - grib_get_process = subprocess.Popen(grib_get_args, stdout=subprocess.PIPE) - uniq_output = subprocess.check_output((uniq_cmd,), stdin=grib_get_process.stdout) - output_paths = [] - skipped_paths = [] - for line in uniq_output.decode('utf-8').rstrip('\n').split('\n'): - splits = dict(zip(split_dims, line.split(' '))) - output_path = self.output_info.formatted_output_path(splits) - if self.should_skip_file(output_path): - skipped_paths.append(output_path) - continue - output_paths.append(output_path) - if not output_paths: - metrics.Metrics.counter('file_splitters', 'skipped').inc() - self.logger.info('Skipping %s, file already split into: %s', - repr(self.input_path), ', '.join(skipped_paths)) - return - - with tempfile.TemporaryDirectory() as tmpdir: - self.logger.info('Performing split.') - dest = os.path.join(tmpdir, flat_output_template) + try: + self.logger.info('Skipping as needed...') + # Append -w flag to filter GRIB messages matching the given expression if self.grib_filter_expression: - subprocess.run([grib_copy_cmd, "-w", - self.grib_filter_expression, - local_file.name, dest], check=True) + grib_get_args = [grib_get_cmd, '-p', split_dims_arg, '-w', self.grib_filter_expression, local_file.name] else: - subprocess.run([grib_copy_cmd, local_file.name, dest], - check=True) - - self.logger.info('Uploading %r...', self.input_path) - for flat_target in os.listdir(tmpdir): - dest_file_path = f'{prefix}{flat_target.replace(delimiter, slash)}' - self.logger.info([prefix, dest_file_path, local_file.name, - self.output_info.unformatted_output_path()]) - - copy(os.path.join(tmpdir, flat_target), dest_file_path) - self.logger.info('Finished uploading %r', self.input_path) - + grib_get_args = [grib_get_cmd, '-p', split_dims_arg, local_file.name] + grib_get_process = subprocess.Popen(grib_get_args, stdout=subprocess.PIPE) + uniq_output = subprocess.check_output((uniq_cmd,), stdin=grib_get_process.stdout) + output_paths = [] + skipped_paths = [] + for line in uniq_output.decode('utf-8').rstrip('\n').split('\n'): + splits = dict(zip(split_dims, line.split(' '))) + output_path = self.output_info.formatted_output_path(splits) + if self.should_skip_file(output_path): + skipped_paths.append(output_path) + continue + output_paths.append(output_path) + if not output_paths: + metrics.Metrics.counter('file_splitters', 'skipped').inc() + self.logger.info('Skipping %s, file already split into: %s', + repr(self.input_path), ', '.join(skipped_paths)) + return + + with tempfile.TemporaryDirectory() as tmpdir: + self.logger.info('Performing split.') + dest = os.path.join(tmpdir, flat_output_template) + if self.grib_filter_expression: + subprocess.run([grib_copy_cmd, "-w", + self.grib_filter_expression, + local_file.name, dest], check=True) + else: + subprocess.run([grib_copy_cmd, local_file.name, dest], + check=True) + + self.logger.info('Uploading %r...', self.input_path) + for flat_target in os.listdir(tmpdir): + dest_file_path = f'{prefix}{flat_target.replace(delimiter, slash)}' + self.logger.info([prefix, dest_file_path, local_file.name, + self.output_info.unformatted_output_path()]) + + copy(os.path.join(tmpdir, flat_target), dest_file_path) + self.logger.info('Finished uploading %r', self.input_path) + except Exception as e: + log_msg = ( + f"GRIB tool failed for {self.input_path!r}. This means the requested " + f"filter expression {self.grib_filter_expression!r} does not exist in this file. " + f"Error: {e}" + ) + if self.skip_on_invalid_grib_filter_expression: + self.logger.warning(f"{log_msg} | Flag 'skip_on_invalid_grib_filter_expression' is True. Skipping file.") + return + else: + self.logger.error(f"{log_msg} | Flag 'skip_on_invalid_grib_filter_expression' is False. Error raised") + raise class NetCdfSplitter(FileSplitter): @@ -353,7 +366,8 @@ def get_splitter(file_path: str, dry_run: bool, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None) -> FileSplitter: + grib_filter_expression: t.Optional[str] = None, + skip_on_invalid_grib_filter_expression: bool = False) -> FileSplitter: if dry_run: logger.info('Using splitter: DrySplitter') return DrySplitter(file_path, output_info, logging_level=logging_level) @@ -371,7 +385,7 @@ def get_splitter(file_path: str, if cmd: logger.info('Using splitter: GribSplitterV2') return GribSplitterV2(file_path, output_info, force_split, - logging_level, grib_filter_expression) + logging_level, grib_filter_expression, skip_on_invalid_grib_filter_expression) else: logger.info('Using splitter: GribSplitter') return GribSplitter(file_path, output_info, force_split, diff --git a/weather_sp/splitter_pipeline/pipeline.py b/weather_sp/splitter_pipeline/pipeline.py index 17b3251..7eb4547 100644 --- a/weather_sp/splitter_pipeline/pipeline.py +++ b/weather_sp/splitter_pipeline/pipeline.py @@ -46,7 +46,8 @@ def split_file(input_file: str, dry_run: bool, force_split: bool = False, logging_level: int = logging.INFO, - grib_filter_expression: t.Optional[str] = None): + grib_filter_expression: t.Optional[str] = None, + skip_on_invalid_grib_filter_expression: bool = False): output_base_name = get_output_base_name(input_path=input_file, input_base=input_base_dir, output_template=output_template, @@ -61,7 +62,8 @@ def split_file(input_file: str, dry_run, force_split, level, - grib_filter_expression) + grib_filter_expression, + skip_on_invalid_grib_filter_expression) splitter.split_data() @@ -134,6 +136,10 @@ def run(argv: t.List[str], save_main_session: bool = True): 'specifically supported by the GribSplitterV2' 'implementation.' 'Example: typeOfLevel=isobaricInhPa,level=1000') + parser.add_argument('--skip-on-invalid-grib-filter-expression', action='store_true', default=False, + help='If provided (True), files that do not contain the key-values specified ' + 'in the --where filter will be logged and skipped. By default (False), ' + 'the pipeline will raise an error and break.') parser.add_argument('--topic', type=str, default=None, help='Pub/Sub topic to read from for streaming mode.') parser.add_argument('--subscription', type=str, default=None, @@ -162,6 +168,7 @@ def run(argv: t.List[str], save_main_session: bool = True): formatting = known_args.formatting dry_run = known_args.dry_run grib_filter_expression = known_args.where + skip_on_invalid_grib_filter_expression = known_args.skip_on_invalid_grib_filter_expression if not output_template and not output_dir: raise ValueError('No output specified') @@ -214,6 +221,7 @@ def run(argv: t.List[str], save_main_session: bool = True): known_args.force, known_args.log_level, grib_filter_expression, + skip_on_invalid_grib_filter_expression, ) )