From 2a3c3fa8bae36adafd99c3c41c2cc4225e00b921 Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Fri, 2 Oct 2026 09:05:44 +0200 Subject: [PATCH 1/3] fix(data-imports): retry object-store timeouts writing parquet batches The pipeline_v3 extraction writer's local retry only matched a bare OSError, so a transient read/connect timeout to object storage skipped the retry and failed the batch on the first attempt. Widen the retry check to also match botocore's ConnectionError/HTTPClientError family, the same classes s3fs itself already treats as retryable internally. Generated-By: PostHog Desktop Task-Id: ac6120a9-5d51-4282-bb2a-69056a55c2f5 --- .../pipelines/pipeline_v3/s3/test_writer.py | 18 ++++++++++++++++++ .../pipelines/pipeline_v3/s3/writer.py | 10 ++++++++++ 2 files changed, 28 insertions(+) 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..c3fb2831bba7 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,23 @@ 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_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..f36fdbe5efb6 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,15 @@ 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). + if isinstance(exc, botocore.exceptions.HTTPClientError | botocore.exceptions.ConnectionError): + return True return type(exc) is OSError From d06e59c16c27d16cdc758a9f9e40e0cdf0c1d48d Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Fri, 2 Oct 2026 09:22:51 +0200 Subject: [PATCH 2/3] fix(data-imports): exclude SSLError from the s3 write retry SSLError is a ConnectionError subclass but usually signals a bad or expired certificate, not a transient blip. Retrying it just delays a failure that happens on every attempt. Addresses a CodeRabbit review comment on #110533. Generated-By: PostHog Desktop Task-Id: ac6120a9-5d51-4282-bb2a-69056a55c2f5 --- .../pipelines/pipeline_v3/s3/test_writer.py | 15 +++++++++++++++ .../pipelines/pipeline_v3/s3/writer.py | 6 +++++- 2 files changed, 20 insertions(+), 1 deletion(-) 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 c3fb2831bba7..5e48155d03ac 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 @@ -68,6 +68,21 @@ def test_retries_read_timeout_then_succeeds(self, mock_write_table, _mock_sleep) 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="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 f36fdbe5efb6..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 @@ -86,7 +86,11 @@ def _is_transient_s3_write_error(exc: BaseException) -> bool: # 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). + # 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 From 010aa15454a72b0e5fce8d47effd44183e3b0a86 Mon Sep 17 00:00:00 2001 From: Tom Owers Date: Fri, 2 Oct 2026 10:34:56 +0200 Subject: [PATCH 3/3] fix(data-imports): construct SSL test error with exception --- .../data_imports/pipelines/pipeline_v3/s3/test_writer.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 5e48155d03ac..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 @@ -75,7 +75,7 @@ def test_does_not_retry_ssl_error(self, mock_write_table) -> None: f = MagicMock() s3 = _fake_s3([f]) mock_write_table.side_effect = botocore.exceptions.SSLError( - endpoint_url="https://example.com/part-0000.parquet", error="certificate verify failed" + endpoint_url="https://example.com/part-0000.parquet", error=Exception("certificate verify failed") ) with pytest.raises(botocore.exceptions.SSLError):