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

Reply via email to