This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 5ec2045d474 Fix PostgresToGCSOperator fetching one row at a time with
psycopg3 (#73324)
5ec2045d474 is described below
commit 5ec2045d4742750d25a73f2ec95f21301d79b561
Author: Shikhar Goel <[email protected]>
AuthorDate: Tue Oct 6 04:05:24 2026 +0530
Fix PostgresToGCSOperator fetching one row at a time with psycopg3 (#73324)
With the psycopg3 driver the server-side cursor path read rows with
fetchone(), which issues FETCH FORWARD 1 per row and ignores
cursor_itersize. Exports of large tables became over an order of
magnitude slower than with psycopg2.
closes: #72075
Co-authored-by: sgoel2be24-cyber
<[email protected]>
---
.../google/cloud/transfers/postgres_to_gcs.py | 14 ++---
.../google/cloud/transfers/test_postgres_to_gcs.py | 61 +++++++++++++++++++++-
2 files changed, 64 insertions(+), 11 deletions(-)
diff --git
a/providers/google/src/airflow/providers/google/cloud/transfers/postgres_to_gcs.py
b/providers/google/src/airflow/providers/google/cloud/transfers/postgres_to_gcs.py
index ff3a7217985..091b0c49298 100644
---
a/providers/google/src/airflow/providers/google/cloud/transfers/postgres_to_gcs.py
+++
b/providers/google/src/airflow/providers/google/cloud/transfers/postgres_to_gcs.py
@@ -49,6 +49,9 @@ class _PostgresServerSideCursorDecorator:
def __init__(self, cursor):
self.cursor = cursor
+ # Iterating the cursor fetches ``itersize`` rows per round trip,
unlike ``fetchone()``. Keep a
+ # single iterator because psycopg < 3.3 returns a new generator on
each ``iter()`` call.
+ self._rows_iterator = iter(cursor)
self.rows = []
self.initialized = False
@@ -58,19 +61,10 @@ class _PostgresServerSideCursorDecorator:
def __next__(self):
"""Fetch next row from the cursor."""
- if USE_PSYCOPG3:
- if self.rows:
- return self.rows.pop()
- self.initialized = True
- row = self.cursor.fetchone()
- if row is None:
- raise StopIteration
- return row
- # psycopg2
if self.rows:
return self.rows.pop()
self.initialized = True
- return next(self.cursor)
+ return next(self._rows_iterator)
@property
def description(self):
diff --git
a/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py
b/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py
index aa0c21ee7e6..5d90b6772d8 100644
--- a/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py
+++ b/providers/google/tests/unit/google/cloud/transfers/test_postgres_to_gcs.py
@@ -28,7 +28,10 @@ from airflow.providers.common.compat.openlineage.facet
import (
SchemaDatasetFacetFields,
)
from airflow.providers.common.sql.hooks.sql import DbApiHook
-from airflow.providers.google.cloud.transfers.postgres_to_gcs import
PostgresToGCSOperator
+from airflow.providers.google.cloud.transfers.postgres_to_gcs import (
+ PostgresToGCSOperator,
+ _PostgresServerSideCursorDecorator,
+)
from airflow.providers.postgres.hooks.postgres import PostgresHook
TABLES = {"postgres_to_gcs_operator", "postgres_to_gcs_operator_empty"}
@@ -53,6 +56,62 @@ SCHEMA_JSON = (
)
+class _FakeServerCursor:
+ """Server-side cursor recording the size of every FETCH; ``__iter__``
pages like psycopg < 3.3."""
+
+ description = [("some_num", 23, None, None, None, None, None)]
+
+ def __init__(self, rows):
+ self._rows = list(rows)
+ self.itersize = 100
+ self.fetch_sizes = []
+
+ def _fetch(self, size):
+ self.fetch_sizes.append(size)
+ batch, self._rows = self._rows[:size], self._rows[size:]
+ return batch
+
+ def fetchone(self):
+ batch = self._fetch(1)
+ return batch[0] if batch else None
+
+ def __iter__(self):
+ while True:
+ batch = self._fetch(self.itersize)
+ yield from batch
+ if len(batch) < self.itersize:
+ return
+
+
+class _FakeSelfIteratingServerCursor(_FakeServerCursor):
+ """Pages like psycopg >= 3.3 and psycopg2 named cursors, which are their
own iterators."""
+
+ _page = None
+ _pos = 0
+
+ def __iter__(self):
+ return self
+
+ def __next__(self):
+ if self._page is None or self._pos >= len(self._page) >= self.itersize:
+ self._page, self._pos = self._fetch(self.itersize), 0
+ if self._pos >= len(self._page):
+ raise StopIteration
+ self._pos += 1
+ return self._page[self._pos - 1]
+
+
[email protected]("cursor_class", [_FakeServerCursor,
_FakeSelfIteratingServerCursor])
+def
test_server_side_cursor_decorator_fetches_in_itersize_batches(cursor_class):
+ rows = [(i,) for i in range(250)]
+ cursor = cursor_class(rows)
+ decorated = _PostgresServerSideCursorDecorator(cursor)
+
+ assert decorated.description == cursor.description
+ assert list(decorated) == rows
+ assert cursor.fetch_sizes == [1, 100, 100, 100]
+
+
@pytest.mark.backend("postgres")
class TestPostgresToGoogleCloudStorageOperator:
@classmethod