Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)).
Expand Down
25 changes: 14 additions & 11 deletions sdks/python/apache_beam/io/fileio.py
Original file line number Diff line number Diff line change
Expand Up @@ -889,17 +889,20 @@ 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)

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

Expand Down
95 changes: 95 additions & 0 deletions sdks/python/apache_beam/io/fileio_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,13 +20,15 @@
# pytype: skip-file

import csv
import errno
import io
import json
import logging
import os
import unittest
import uuid
import warnings
from unittest import mock

import pytest
from hamcrest.library.text import stringmatches
Expand All @@ -40,6 +42,7 @@
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
Expand Down Expand Up @@ -766,6 +769,98 @@ 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 = [
Expand Down
Loading