diff --git a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/test_writer.py b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/test_writer.py index 265f6c703151..d3bdab851ae5 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/test_writer.py +++ b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/test_writer.py @@ -6,6 +6,7 @@ from unittest.mock import MagicMock, patch import pyarrow as pa +import botocore.exceptions from parameterized import parameterized from products.warehouse_sources.backend.temporal.data_imports.pipelines.pipeline_v3.s3.writer import ( @@ -50,6 +51,38 @@ def test_reraises_after_persistent_os_error(self, mock_write_table, _mock_sleep) assert mock_write_table.call_count == 4 + @patch("tenacity.nap.time.sleep") + @patch("pyarrow.parquet.write_table") + def test_retries_read_timeout_then_succeeds(self, mock_write_table, _mock_sleep) -> None: + # A read/connect timeout never gets s3fs's OSError translation (there's no response to + # translate), so it reaches here as the raw botocore exception rather than an OSError. + # It's exactly the transient blip this retry exists for and shouldn't fail the batch. + f = MagicMock() + s3 = _fake_s3([f, f]) + mock_write_table.side_effect = [ + botocore.exceptions.ReadTimeoutError(endpoint_url="https://example.com/part-0000.parquet"), + None, + ] + + _write_parquet_to_s3(s3, "bucket/part-0000.parquet", pa.table({"id": [1]}), "zstd") + + assert mock_write_table.call_count == 2 + + @patch("pyarrow.parquet.write_table") + def test_does_not_retry_ssl_error(self, mock_write_table) -> None: + # SSLError is a ConnectionError subclass but usually means a bad/expired certificate, + # not a network blip — retrying just delays a failure that will happen on every attempt. + f = MagicMock() + s3 = _fake_s3([f]) + mock_write_table.side_effect = botocore.exceptions.SSLError( + endpoint_url="https://example.com/part-0000.parquet", error=Exception("certificate verify failed") + ) + + with pytest.raises(botocore.exceptions.SSLError): + _write_parquet_to_s3(s3, "bucket/part-0000.parquet", pa.table({"id": [1]}), "zstd") + + assert mock_write_table.call_count == 1 + @patch("pyarrow.parquet.write_table") def test_does_not_retry_permission_error(self, mock_write_table) -> None: # PermissionError (e.g. AccessDenied/InvalidAccessKeyId) is an OSError subclass but not diff --git a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/writer.py b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/writer.py index f03473f2be83..49070a34e28e 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/writer.py +++ b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/s3/writer.py @@ -6,6 +6,7 @@ import s3fs import pyarrow as pa import pyarrow.parquet as pq +import botocore.exceptions from structlog.types import FilteringBoundLogger from temporalio import activity from tenacity import retry, retry_if_exception, stop_after_attempt, wait_exponential_jitter @@ -79,6 +80,19 @@ def _is_transient_s3_write_error(exc: BaseException) -> bool: # FileNotFoundError, and TimeoutError are also OSError subclasses but signal # non-transient causes (bad credentials, deleted bucket) that retrying won't fix, # so only the exact base type is treated as retryable here. + # + # A read/connect timeout or a dropped connection to the object store never gets that OSError + # translation: s3fs's own internal retries only convert a *response* it got back into an + # OSError, and a timed-out or refused connection never got one. It reaches here as the raw + # botocore exception once those internal retries are exhausted, so it needs a type check of + # its own (these are exactly the classes s3fs itself treats as retryable, see its + # S3_RETRYABLE_ERRORS/ClientError handling in s3fs.core._error_wrapper). SSLError is a + # ConnectionError subclass but usually means a persistent certificate problem, not a blip, + # so it's excluded rather than spending the whole retry budget before failing anyway. + if isinstance(exc, botocore.exceptions.SSLError): + return False + if isinstance(exc, botocore.exceptions.HTTPClientError | botocore.exceptions.ConnectionError): + return True return type(exc) is OSError