This is an automated email from the ASF dual-hosted git repository.
jrmccluskey 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 f5ba5fa5910 [Python] Propagate WriteToFiles finalization failures
(#39993)
f5ba5fa5910 is described below
commit f5ba5fa5910dca23f4b150ec11cee94c081b1a43
Author: Bruno Volpato <[email protected]>
AuthorDate: Thu Sep 10 10:12:14 2026 -0400
[Python] Propagate WriteToFiles finalization failures (#39993)
* [Python] Propagate WriteToFiles finalization failures
* [Python] Link file finalization fix in changelog
---
CHANGES.md | 1 +
sdks/python/apache_beam/io/fileio.py | 25 ++++----
sdks/python/apache_beam/io/fileio_test.py | 95 +++++++++++++++++++++++++++++++
3 files changed, 110 insertions(+), 11 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index f0d5d06b9d9..e58f73c9f2a 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -93,6 +93,7 @@
## 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 fcce83fa59e..a51c272d6c7 100644
--- a/sdks/python/apache_beam/io/fileio.py
+++ b/sdks/python/apache_beam/io/fileio.py
@@ -889,17 +889,20 @@ 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)
- 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)
+ 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)
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 1c650653804..87215ac3beb 100644
--- a/sdks/python/apache_beam/io/fileio_test.py
+++ b/sdks/python/apache_beam/io/fileio_test.py
@@ -20,6 +20,7 @@
# pytype: skip-file
import csv
+import errno
import io
import json
import logging
@@ -27,6 +28,7 @@ import os
import unittest
import uuid
import warnings
+from unittest import mock
import pytest
from hamcrest.library.text import stringmatches
@@ -40,6 +42,7 @@ 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
@@ -766,6 +769,98 @@ 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 = [