diff --git a/CHANGES.md b/CHANGES.md index e58f73c9f2a0..f0d5d06b9d9f 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -93,7 +93,6 @@ ## Bugfixes * (Java) Fixed the Spark runner firing processing-time timers in reverse timestamp order ([#39824](https://github.com/apache/beam/issues/39824)). -* (Python) `WriteToFiles` now propagates failed final file moves while preserving retries of already completed moves ([#39993](https://github.com/apache/beam/pull/39993)). * (Python) Fixed incorrect profiler options handling on portable runners ([#39613](https://github.com/apache/beam/issues/39613)). * (Java) KafkaIO dynamic reads no longer require the obsolete `beam_fn_api` experiment ([#29998](https://github.com/apache/beam/issues/29998)). * (Prism) Self-checkpointing splittable DoFns now resume after their requested delay instead of immediately, so polling SDFs no longer busy-spin ([#39848](https://github.com/apache/beam/issues/39848)). diff --git a/sdks/python/apache_beam/io/fileio.py b/sdks/python/apache_beam/io/fileio.py index a51c272d6c75..fcce83fa59ec 100644 --- a/sdks/python/apache_beam/io/fileio.py +++ b/sdks/python/apache_beam/io/fileio.py @@ -889,20 +889,17 @@ def process(self, element, w=beam.DoFn.WindowParam): # Usually harmless. Especially if see FileExistsError so no need to log _LOGGER.debug('Fail to create dir for final destination: %s', cause) - pending_sources = [] - pending_destinations = [] - for source, name in zip(move_from, move_to): - target = filesystems.FileSystems.join(self.path.get(), name) - # A previous bundle attempt may have already moved some or all files. - # An existing target alone is insufficient: it may contain older data. - if (not filesystems.FileSystems.exists(source) and - filesystems.FileSystems.exists(target)): - continue - pending_sources.append(source) - pending_destinations.append(target) - - if pending_sources: - filesystems.FileSystems.rename(pending_sources, pending_destinations) + try: + filesystems.FileSystems.rename( + move_from, + [filesystems.FileSystems.join(self.path.get(), f) for f in move_to]) + except BeamIOError: + # This error is not serious, because it may happen on a retry of the + # bundle. We simply log it. + _LOGGER.debug( + 'Exception occurred during moving files: %s. This may be due to a' + ' bundle being retried.', + move_from) yield from final_file_results diff --git a/sdks/python/apache_beam/io/fileio_test.py b/sdks/python/apache_beam/io/fileio_test.py index 87215ac3beb3..1c650653804c 100644 --- a/sdks/python/apache_beam/io/fileio_test.py +++ b/sdks/python/apache_beam/io/fileio_test.py @@ -20,7 +20,6 @@ # pytype: skip-file import csv -import errno import io import json import logging @@ -28,7 +27,6 @@ import unittest import uuid import warnings -from unittest import mock import pytest from hamcrest.library.text import stringmatches @@ -42,7 +40,6 @@ from apache_beam.io.filesystems import FileSystems from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.options.pipeline_options import StandardOptions -from apache_beam.options.value_provider import StaticValueProvider from apache_beam.testing.test_pipeline import TestPipeline from apache_beam.testing.test_stream import TestStream from apache_beam.testing.test_utils import compute_hash @@ -769,98 +766,6 @@ def _touch(element): assert_that(match_continiously, equal_to([path, path])) -class MoveTempFilesTest(_TestCaseWithTempDirCleanUp): - def setUp(self): - super().setUp() - self.temp_dir = self._new_tempdir() - self.output_dir = self._new_tempdir() - self.move_files = fileio._MoveTempFilesIntoFinalDestinationFn( - StaticValueProvider(str, self.output_dir), lambda window, pane, shard, - total, compression, destination: 'part-%d' % shard, - StaticValueProvider(str, self.temp_dir)) - - def _file_results(self, *contents): - return [ - fileio.FileResult( - self._create_temp_file(dir=self.temp_dir, content=content), - shard_index=i, - total_shards=len(contents), - window=GlobalWindow(), - pane=None, - destination='destination') for i, content in enumerate(contents) - ] - - def _finalize(self, results): - return self.move_files.process(('destination', results), w=GlobalWindow()) - - def _output_path(self, shard): - return FileSystems.join(self.output_dir, 'part-%d' % shard) - - def test_rename_failure_does_not_emit_results(self): - results = self._file_results('row') - error = BeamIOError('Rename failed') - with mock.patch.object(FileSystems, 'rename', side_effect=error): - with self.assertRaises(BeamIOError) as raised: - next(self._finalize(results)) - self.assertIs(raised.exception, error) - self.assertTrue(FileSystems.exists(results[0].file_name)) - self.assertFalse(FileSystems.exists(self._output_path(0))) - - def test_completed_rename_can_be_retried(self): - results = self._file_results('row') - expected = list(self._finalize(results)) - self.assertEqual(expected, list(self._finalize(results))) - with open(self._output_path(0)) as output: - self.assertEqual('row', output.read()) - - def test_partial_rename_failure_can_be_retried(self): - results = self._file_results('first row', 'second row') - original_rename = os.rename - - def fail_second_file(source, destination): - if source == results[1].file_name: - raise OSError(errno.EIO, 'Temporary I/O error') - original_rename(source, destination) - - with mock.patch('os.rename', side_effect=fail_second_file): - with self.assertRaises(BeamIOError): - next(self._finalize(results)) - self.assertFalse(FileSystems.exists(results[0].file_name)) - self.assertTrue(FileSystems.exists(results[1].file_name)) - - with mock.patch.object(FileSystems, 'rename', - wraps=FileSystems.rename) as rename: - self.assertEqual(2, len(list(self._finalize(results)))) - rename.assert_called_once_with([results[1].file_name], - [self._output_path(1)]) - for shard, expected in enumerate(['first row', 'second row']): - with open(self._output_path(shard)) as output: - self.assertEqual(expected, output.read()) - - def test_existing_destination_does_not_hide_rename_failure(self): - results = self._file_results('new row') - target = self._output_path(0) - with open(target, 'w') as output: - output.write('old row') - error = BeamIOError( - 'Rename failed', - {(results[0].file_name, target): PermissionError('Access denied')}) - with mock.patch.object(FileSystems, 'rename', side_effect=error): - with self.assertRaises(BeamIOError) as raised: - next(self._finalize(results)) - self.assertIs(raised.exception, error) - with open(results[0].file_name) as source: - self.assertEqual('new row', source.read()) - with open(target) as output: - self.assertEqual('old row', output.read()) - - def test_missing_source_and_destination_fails(self): - results = self._file_results('row') - FileSystems.delete([results[0].file_name]) - with self.assertRaises(BeamIOError): - next(self._finalize(results)) - - class WriteFilesTest(_TestCaseWithTempDirCleanUp): SIMPLE_COLLECTION = [