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 = [

Reply via email to