This is an automated email from the ASF dual-hosted git repository.
Abacn 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 634bfe71a97 Revert "[Python] Propagate WriteToFiles finalization
failures (#39993)" (#40095)
634bfe71a97 is described below
commit 634bfe71a97747c1728e4fa1620dceb5edc94eed
Author: Yi Hu <[email protected]>
AuthorDate: Thu Sep 10 14:04:58 2026 -0400
Revert "[Python] Propagate WriteToFiles finalization failures (#39993)"
(#40095)
---
CHANGES.md | 1 -
sdks/python/apache_beam/io/fileio.py | 25 ++++----
sdks/python/apache_beam/io/fileio_test.py | 95 -------------------------------
3 files changed, 11 insertions(+), 110 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index e58f73c9f2a..f0d5d06b9d9 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 a51c272d6c7..fcce83fa59e 100644
--- a/sdks/python/apache_beam/io/fileio.py
+++ b/sdks/python/apache_beam/io/fileio.py
@@ -889,20 +889,17 @@ class _MoveTempFilesIntoFinalDestinationFn(beam.DoFn):
# 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 87215ac3beb..1c650653804 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 os
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.filesystem import FileMetadata
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 @@ class MatchContinuouslyTest(_TestCaseWithTempDirCleanUp):
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 = [