Skip to content

Fix PostgresToGCSOperator fetching one row at a time with psycopg3 - #73324

Merged
potiuk merged 1 commit into
apache:mainfrom
sgoel2be24-cyber:fix-postgres-to-gcs-psycopg3-itersize
Oct 5, 2026
Merged

potiuk merged 1 commit into
apache:mainfrom
sgoel2be24-cyber:fix-postgres-to-gcs-psycopg3-itersize

Conversation

@sgoel2be24-cyber

Copy link
Copy Markdown
Contributor

With the psycopg3 driver (the default since apache-airflow-providers-postgres 7.0.0), PostgresToGCSOperator(use_server_side_cursor=True) read rows through fetchone(), which issues FETCH FORWARD 1 per row and ignores cursor_itersize. The issue reports a 2M-row export going from ~5 to 90+ minutes.

_PostgresServerSideCursorDecorator now iterates the cursor for both drivers, so psycopg fetches itersize rows per round trip, as psycopg2 already did. The decorator keeps a single iterator because psycopg < 3.3 (the provider allows >=3.2.9) implements ServerCursor.__iter__ as a generator: calling iter() again would start a new one and drop the rest of the current batch. psycopg >= 3.3 and psycopg2 named cursors are their own iterators, so nothing changes for psycopg2.

The existing tests for this operator need a Postgres backend, and they pass with any fetch size, so they couldn't catch this. The new unit test uses fake cursors that behave like both psycopg iteration styles and asserts every row is returned and rows are fetched in itersize batches ([1, 100, 100, 100] for 250 rows; the single-row fetch is the one description needs). On main both cases fail with 251 single-row fetches.

Tested locally:

  • pytest on test_postgres_to_gcs.py and test_sql_to_gcs.py: 16 passed. The 18 Postgres-backend tests were skipped: there's no Postgres or Docker on my machine, so they'll run in CI
  • prek pre-commit stage passes (check-template-fields-valid needs Docker and was skipped; no template fields changed); mypy on postgres_to_gcs.py passes
  • Checked the ServerCursor source for psycopg 3.2.9 and 3.3.5 to confirm both iteration behaviours the fakes model

closes: #72075


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (Opus 5)

Generated-by: Claude Code (Opus 5) following the guidelines

@boring-cyborg boring-cyborg Bot added area:providers provider:google Google (including GCP) related issues labels Sep 18, 2026
@sgoel2be24-cyber

sgoel2be24-cyber commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor Author

Gentle ping for review when someone has a moment. This fixes #72075, where psycopg3's server-side cursor was fetching one row at a time. CI is green.


Drafted-by: Claude Code (Opus 5); reviewed by @sgoel2be24-cyber before posting

@molcay

molcay commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

Hi @sgoel2be24-cyber,

Can you share the system tests run results for the PostgresToGCSOperator?

@sgoel2be24-cyber

Copy link
Copy Markdown
Contributor Author

Hi @molcay, thanks for taking a look. I can't run example_postgres_to_gcs myself: it provisions a Compute Engine VM and a GCS bucket, and I don't have a GCP environment set up for it. Instead, I ran the operator end-to-end against a real PostgreSQL server with only the GCS upload stubbed, and counted the FETCH statements Postgres received (log_statement = 'all').

Setup: PostgreSQL 16.2, Airflow 3.3.2, apache-airflow-providers-postgres 7.1.0 (psycopg3 path), and apache-airflow-providers-google 22.6.0, whose postgres_to_gcs.py is identical to current main. I compared that file with the same file plus this PR's patch, using use_server_side_cursor=True and JSON export.

psycopg rows cursor_itersize version FETCH statements received by Postgres rows exported
3.2.9 10,000 1,000 main 10,001 × FETCH FORWARD 1 10,000
3.2.9 10,000 1,000 this PR 1 × FETCH FORWARD 1 + 10 × FETCH FORWARD 1000 10,000
3.3.6 100,000 2,000 main 100,001 × FETCH FORWARD 1 100,000
3.3.6 100,000 2,000 this PR 1 × FETCH FORWARD 1 + 50 × FETCH FORWARD 2000 100,000

In every run the exported rows were complete and in order. The remaining FETCH FORWARD 1 comes from the description property, which reads the first row to populate cursor.description. On localhost, 100k rows went from 6.0s to 3.3s. With real latency between the worker and the database, as in #72075, the gap is much larger because main makes one round trip per row.

I'm happy to share the script. If someone with GCP access could run the system test, I'd appreciate it. Otherwise, let me know if you'd like this verified another way.


Drafted-by: Claude Code (Opus 5.5)

@shahar1
shahar1 requested a review from Dev-iL October 2, 2026 08:01
@Dev-iL

Dev-iL commented Oct 2, 2026

Copy link
Copy Markdown
Collaborator

Thanks for looking into this! I observed something similar recently when migrating internal endpoints from sync to async, and with async operators it became apparent that something was off.

@Dev-iL Dev-iL left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Verified the claims on my end - everything works as advertised. Thank you for your contribution!

Verification performed

Ran a temporary probe through Breeze with its PostgreSQL backend (configured PostgreSQL 14), psycopg 3.3.6, and psycopg2 2.9.13. The probe extracts the exact decorator definitions from base and head, avoiding dependence on the current checkout's operator implementation. It checks 0, 1, 1,000, and 5,000 rows at itersize 1, 100, and 2,000, with four repetitions per combination. All 192 real-driver cases returned complete, ordered rows and the expected description after repeated description access.

The probe also executes the added regression test's definitions against both revisions. On the psycopg3 path, both fake-cursor styles fail the batch assertion on base and pass on head. Both styles pass on head regardless of the driver flag. The artificial generator fake is not a psycopg2 cursor; its base psycopg2 TypeError is not a product defect.

Illustrative warm cursor timings at itersize=2000, excluding the first repetition:

Driver Rows Base median / maximum Head median / maximum
psycopg 3.3.6 1,000 140.8 / 148.5 ms 2.1 / 2.2 ms
psycopg 3.3.6 5,000 712.3 / 783.7 ms 10.6 / 12.3 ms
psycopg2 2.9.13 1,000 1.2 / 1.3 ms 1.3 / 1.4 ms
psycopg2 2.9.13 5,000 6.9 / 7.4 ms 9.1 / 10.9 ms

These are one local run with three warm samples per combination, covering execute, schema access, and cursor consumption. They establish neither production throughput nor a reliable small psycopg2 timing regression. They exclude conversion, file writing, and GCS upload.

@molcay

molcay commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

Hi @Dev-iL,

Is there any chance while you are profiling the operations, you used the real GCS endpoints?
I just wonder if this change is OK e2e.

I did not find a chance to run the system test(s) for this operator and @sgoel2be24-cyber does not have access to GCP.

@Dev-iL

Dev-iL commented Oct 5, 2026

Copy link
Copy Markdown
Collaborator

@molcay I'm afraid I don't have access to a live instance to test against.

@shahar1 Any idea if this is something that might be covered by Google's CI?

@molcay

molcay commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

Hi @Dev-iL,

Thank you for the answer.

I see. I will try to run it locally to check if it is ok or not.

About the Google's CI, currently we are not running against each PR. It is running daily on main branch and also for RCs.

@Dev-iL

Dev-iL commented Oct 5, 2026

Copy link
Copy Markdown
Collaborator

Oh I haven't noticed you're affiliated with google.

I should've clarified that my own tests were focused on the postgres side of things.

Worst case if this is merged and broken, it will be picked up by the next daily run and we can revert or exclude it from the upcoming provider wave, no?

@molcay

molcay commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

No problem :)

Thanks for clarification.


Worst case if this is merged and broken, it will be picked up by the next daily run and we can revert or exclude it from the upcoming provider wave, no?

Yes, but we prefer to catch this kind of things before merging. Today, I will try to execute and check the result

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: apache#72075
@potiuk
potiuk force-pushed the fix-postgres-to-gcs-psycopg3-itersize branch from 4207c5f to 41681ce Compare October 5, 2026 13:39
@potiuk

potiuk commented Oct 5, 2026

Copy link
Copy Markdown
Member

Yeah. Would be great if we can confirm it before merging. Note that I am going to start new provider's release tomorrow, so we might choose to merge it without system tests check @molcay if it's not complete.

@molcay

molcay commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

Hi @potiuk,

I finally completed the test. The change seems OK. We can merge it for tomorrow RC.

@potiuk
potiuk merged commit 5ec2045 into apache:main Oct 5, 2026
87 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providers provider:google Google (including GCP) related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

PostgresToGCSOperator server-side cursor ignores cursor_itersize with psycopg3 driver (FETCH FORWARD 1)

4 participants