This is an automated email from the ASF dual-hosted git repository.
shahar1 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 7c272f1ea97 Make psycopg (v3) the default synchronous Postgres driver
(#69526)
7c272f1ea97 is described below
commit 7c272f1ea97c9052d0ab96a696064d62dea584c1
Author: Dev-iL <[email protected]>
AuthorDate: Wed Jul 22 20:08:50 2026 +0300
Make psycopg (v3) the default synchronous Postgres driver (#69526)
---
docs/spelling_wordlist.txt | 1 +
providers/amazon/docs/changelog.rst | 7 ++
providers/amazon/docs/executors/batch-executor.rst | 2 +-
providers/amazon/docs/executors/ecs-executor.rst | 2 +-
.../amazon/docs/executors/lambda-executor.rst | 2 +-
.../providers/amazon/aws/hooks/redshift_sql.py | 37 ++++++++-
.../unit/amazon/aws/hooks/test_redshift_sql.py | 43 +++++++++-
providers/celery/docs/changelog.rst | 7 ++
providers/celery/provider.yaml | 2 +-
.../providers/celery/executors/default_celery.py | 21 ++++-
.../airflow/providers/celery/get_provider_info.py | 2 +-
.../unit/celery/executors/test_celery_executor.py | 26 ++++++
.../src/airflow/providers/common/sql/hooks/sql.py | 37 ++++++++-
.../sql/tests/unit/common/sql/hooks/test_sql.py | 92 ++++++++++++++++++++++
.../google/cloud/transfers/bigquery_to_postgres.py | 13 +--
.../cloud/transfers/test_bigquery_to_postgres.py | 51 +++++++++++-
providers/pgvector/pyproject.toml | 1 +
.../providers/pgvector/operators/pgvector.py | 21 ++++-
.../tests/unit/pgvector/operators/test_pgvector.py | 57 ++++++++++++--
providers/postgres/README.rst | 5 +-
providers/postgres/docs/changelog.rst | 14 ++++
providers/postgres/docs/index.rst | 5 +-
providers/postgres/pyproject.toml | 24 +++---
.../airflow/providers/postgres/hooks/postgres.py | 42 ++++++++--
.../tests/unit/postgres/hooks/test_postgres.py | 21 +++++
.../unit/postgres/test_provider_dependencies.py | 46 +++++++++++
uv.lock | 9 ++-
27 files changed, 542 insertions(+), 48 deletions(-)
diff --git a/docs/spelling_wordlist.txt b/docs/spelling_wordlist.txt
index 17fa3bc2deb..cf74e1f5f07 100644
--- a/docs/spelling_wordlist.txt
+++ b/docs/spelling_wordlist.txt
@@ -524,6 +524,7 @@ downstreams
dp
dq
Drillbit
+driverless
dropdown
druidHook
ds
diff --git a/providers/amazon/docs/changelog.rst
b/providers/amazon/docs/changelog.rst
index 4f59d4bede4..0e8a27715b6 100644
--- a/providers/amazon/docs/changelog.rst
+++ b/providers/amazon/docs/changelog.rst
@@ -26,6 +26,13 @@
Changelog
---------
+.. note::
+ ``RedshiftSQLHook.get_sqlalchemy_engine()``, used for the SQLAlchemy
engine backing
+ lineage/reflection, now resolves an explicit Postgres DB-API driver for a
bare
+ ``postgresql://`` URL, preferring ``psycopg`` (v3) when importable and
falling back to
+ ``psycopg2``. Where both drivers are installed, this engine now uses
``psycopg`` (v3),
+ matching the sync Postgres driver default change in
``apache-airflow-providers-postgres``.
+
9.32.0
......
diff --git a/providers/amazon/docs/executors/batch-executor.rst
b/providers/amazon/docs/executors/batch-executor.rst
index 64ac98672ec..8f71ba7d8c5 100644
--- a/providers/amazon/docs/executors/batch-executor.rst
+++ b/providers/amazon/docs/executors/batch-executor.rst
@@ -298,7 +298,7 @@ Create a Job Definition
.. code-block:: bash
- postgresql+psycopg2://<username>:<password>@<endpoint>/<database_name>
+ postgresql+psycopg://<username>:<password>@<endpoint>/<database_name>
7. Add other configuration as necessary for Airflow generally (see `here
<https://airflow.apache.org/docs/apache-airflow/stable/configurations-ref.html>`__),
the Batch executor (see :ref:`here <config-options>`) or for remote logging
(see :ref:`here <logging>`). Note that any configuration changes should be made
across the entire Airflow environment to keep configuration consistent.
diff --git a/providers/amazon/docs/executors/ecs-executor.rst
b/providers/amazon/docs/executors/ecs-executor.rst
index 9b98e3fc79a..fd534cb27de 100644
--- a/providers/amazon/docs/executors/ecs-executor.rst
+++ b/providers/amazon/docs/executors/ecs-executor.rst
@@ -332,7 +332,7 @@ Create Task Definition
.. code-block:: bash
- postgresql+psycopg2://<username>:<password>@<endpoint>/<database_name>
+ postgresql+psycopg://<username>:<password>@<endpoint>/<database_name>
- ``AIRFLOW__ECS_EXECUTOR__SECURITY_GROUPS``, with the value being a comma
separated list of security group IDs associated with the VPC used for the RDS
instance.
diff --git a/providers/amazon/docs/executors/lambda-executor.rst
b/providers/amazon/docs/executors/lambda-executor.rst
index 3451ba774ce..4ceb5c3bea9 100644
--- a/providers/amazon/docs/executors/lambda-executor.rst
+++ b/providers/amazon/docs/executors/lambda-executor.rst
@@ -334,7 +334,7 @@ Finally create the function:
.. code-block:: bash
- postgresql+psycopg2://<username>:<password>@<endpoint>/<database_name>
+ postgresql+psycopg://<username>:<password>@<endpoint>/<database_name>
- ``AIRFLOW__LAMBDA_EXECUTOR__QUEUE_URL``, with the value being the URL of the
SQS queue created above.
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/hooks/redshift_sql.py
b/providers/amazon/src/airflow/providers/amazon/aws/hooks/redshift_sql.py
index d3ce5819042..e4af4e71724 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/hooks/redshift_sql.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/hooks/redshift_sql.py
@@ -17,6 +17,7 @@
from __future__ import annotations
from functools import cached_property
+from importlib.util import find_spec
from typing import TYPE_CHECKING
import redshift_connector
@@ -25,9 +26,9 @@ from redshift_connector import Connection as
RedshiftConnection, InterfaceError,
try:
from sqlalchemy import create_engine
- from sqlalchemy.engine.url import URL
+ from sqlalchemy.engine.url import URL, make_url
except ImportError:
- URL = create_engine = None # type: ignore[assignment,misc]
+ URL = create_engine = make_url = None # type: ignore[assignment,misc]
from airflow.providers.amazon.aws.hooks.base_aws import AwsBaseHook
from airflow.providers.common.compat.sdk import AirflowException,
AirflowOptionalProviderFeatureException
@@ -38,6 +39,35 @@ if TYPE_CHECKING:
from airflow.providers.openlineage.sqlparser import DatabaseInfo
+def _is_sqlalchemy_2() -> bool:
+ """Whether the installed SQLAlchemy has the native psycopg (v3) dialect
(added in 2.0)."""
+ from packaging.version import Version
+ from sqlalchemy import __version__ as sqlalchemy_version
+
+ return Version(sqlalchemy_version) >= Version("2.0.0")
+
+
+def _resolve_postgres_drivername() -> str:
+ """
+ Pick a Postgres DB-API driver for the SQLAlchemy engine used by
lineage/reflection.
+
+ Redshift is wire-compatible with Postgres, but Airflow's own connection to
Redshift
+ uses ``redshift_connector`` directly (see ``get_conn``) — this only picks
the driver
+ for the separate SQLAlchemy engine used e.g. for OpenLineage extraction.
+ """
+ # The psycopg (v3) package may be importable even where SQLAlchemy is
pinned below 2.0
+ # (e.g. compat tests against older released Airflow versions), which
doesn't register
+ # SQLAlchemy's native "postgresql+psycopg" dialect and would raise
NoSuchModuleError.
+ if find_spec("psycopg") is not None and _is_sqlalchemy_2():
+ return "postgresql+psycopg"
+ if find_spec("psycopg2") is not None:
+ return "postgresql+psycopg2"
+ raise AirflowOptionalProviderFeatureException(
+ "A Postgres DB-API driver is required for the SQLAlchemy engine.
Install one with: "
+ "pip install 'apache-airflow-providers-amazon[sqlalchemy]'
psycopg2-binary (or psycopg[binary])"
+ )
+
+
class RedshiftSQLHook(DbApiHook):
"""
Execute statements against Amazon Redshift.
@@ -201,7 +231,8 @@ class RedshiftSQLHook(DbApiHook):
else:
engine_kwargs["connect_args"] = conn_kwargs
- return create_engine(self.get_uri(), **engine_kwargs)
+ engine_url =
make_url(self.get_uri()).set(drivername=_resolve_postgres_drivername())
+ return create_engine(engine_url, **engine_kwargs)
def get_table_primary_key(self, table: str, schema: str | None = "public")
-> list[str] | None:
"""
diff --git a/providers/amazon/tests/unit/amazon/aws/hooks/test_redshift_sql.py
b/providers/amazon/tests/unit/amazon/aws/hooks/test_redshift_sql.py
index 80a5877e495..42a2dd36e44 100644
--- a/providers/amazon/tests/unit/amazon/aws/hooks/test_redshift_sql.py
+++ b/providers/amazon/tests/unit/amazon/aws/hooks/test_redshift_sql.py
@@ -24,7 +24,7 @@ import pytest
from airflow.models import Connection
from airflow.providers.amazon.aws.hooks.redshift_sql import RedshiftSQLHook
from airflow.providers.amazon.version_compat import NOTSET
-from airflow.providers.common.compat.sdk import AirflowException
+from airflow.providers.common.compat.sdk import AirflowException,
AirflowOptionalProviderFeatureException
LOGIN_USER = "login"
LOGIN_PASSWORD = "password"
@@ -54,6 +54,47 @@ class TestRedshiftSQLHookConn:
expected = "postgresql://login:password@host:5439/dev"
assert db_uri == expected
+
@mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.create_engine")
+
@mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql._is_sqlalchemy_2",
return_value=True)
+ @mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.find_spec")
+ def test_get_sqlalchemy_engine_prefers_psycopg3(self, mock_find_spec,
mock_is_sqla2, mock_create_engine):
+ # Mock create_engine so this doesn't depend on psycopg actually being
importable, or on
+ # the ambient SQLAlchemy major version matching what _is_sqlalchemy_2
is mocked to return.
+ mock_find_spec.return_value = object()
+ self.db_hook.get_sqlalchemy_engine()
+ engine_url = mock_create_engine.call_args[0][0]
+ assert engine_url.drivername == "postgresql+psycopg"
+
+
@mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.create_engine")
+ @mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.find_spec")
+ def test_get_sqlalchemy_engine_falls_back_to_psycopg2(self,
mock_find_spec, mock_create_engine):
+ # psycopg2 is optional (see providers/postgres's [psycopg2] extra) and
may not be
+ # installed here, so mock create_engine instead of letting it actually
import the driver.
+ mock_find_spec.side_effect = lambda name: None if name == "psycopg"
else object()
+ self.db_hook.get_sqlalchemy_engine()
+ engine_url = mock_create_engine.call_args[0][0]
+ assert engine_url.drivername == "postgresql+psycopg2"
+
+
@mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.create_engine")
+
@mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql._is_sqlalchemy_2",
return_value=False)
+ @mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.find_spec")
+ def test_get_sqlalchemy_engine_falls_back_to_psycopg2_on_sqlalchemy_1(
+ self, mock_find_spec, mock_is_sqla2, mock_create_engine
+ ):
+ # psycopg (v3) can be importable even when SQLAlchemy is pinned below
2.0 (e.g. compat
+ # tests against older released Airflow versions) — SQLAlchemy's native
"postgresql+psycopg"
+ # dialect only exists from 2.0 onwards, so this must still fall back
to psycopg2 rather
+ # than let create_engine raise NoSuchModuleError.
+ mock_find_spec.return_value = object()
+ self.db_hook.get_sqlalchemy_engine()
+ engine_url = mock_create_engine.call_args[0][0]
+ assert engine_url.drivername == "postgresql+psycopg2"
+
+ @mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.find_spec",
return_value=None)
+ def test_get_sqlalchemy_engine_raises_clear_error_without_a_driver(self,
mock_find_spec):
+ with pytest.raises(AirflowOptionalProviderFeatureException,
match="Postgres DB-API driver"):
+ self.db_hook.get_sqlalchemy_engine()
+
@mock.patch("airflow.providers.amazon.aws.hooks.redshift_sql.redshift_connector.connect")
def test_get_conn(self, mock_connect):
self.db_hook.get_conn()
diff --git a/providers/celery/docs/changelog.rst
b/providers/celery/docs/changelog.rst
index 1c3683bd63f..fba28081d8c 100644
--- a/providers/celery/docs/changelog.rst
+++ b/providers/celery/docs/changelog.rst
@@ -27,6 +27,13 @@
Changelog
---------
+.. note::
+ When ``[celery] result_backend`` is not set and Airflow derives it from
``sql_alchemy_conn``, a
+ driverless ``postgresql://`` connection string now produces
``db+postgresql+psycopg://...`` instead
+ of ``db+postgresql+psycopg2://...``, matching the sync Postgres driver
default change in
+ ``apache-airflow`` and ``apache-airflow-providers-postgres``. An explicit
``[celery] result_backend``
+ or an explicit driver in ``sql_alchemy_conn`` is unaffected.
+
3.22.0
......
diff --git a/providers/celery/provider.yaml b/providers/celery/provider.yaml
index 1e52e95201b..f5506e41122 100644
--- a/providers/celery/provider.yaml
+++ b/providers/celery/provider.yaml
@@ -248,7 +248,7 @@ config:
version_added: ~
type: string
sensitive: true
- example: "db+postgresql+psycopg2://postgres:airflow@postgres/airflow"
+ example: "db+postgresql+psycopg://postgres:airflow@postgres/airflow"
default: ~
result_backend_sqlalchemy_engine_options:
description: |
diff --git
a/providers/celery/src/airflow/providers/celery/executors/default_celery.py
b/providers/celery/src/airflow/providers/celery/executors/default_celery.py
index 9dd94036734..3af3aa7744e 100644
--- a/providers/celery/src/airflow/providers/celery/executors/default_celery.py
+++ b/providers/celery/src/airflow/providers/celery/executors/default_celery.py
@@ -31,6 +31,18 @@ from airflow.providers.common.compat.sdk import
AirflowException, conf
log = logging.getLogger(__name__)
+_USE_PSYCOPG3: bool
+try:
+ from importlib.metadata import version
+ from importlib.util import find_spec
+
+ from packaging.version import Version
+
+ _is_sqla2 = Version(version("sqlalchemy")) >= Version("2.0.0")
+ _USE_PSYCOPG3 = _is_sqla2 and find_spec("psycopg") is not None
+except (ImportError, ModuleNotFoundError):
+ _USE_PSYCOPG3 = False
+
def _broker_supports_visibility_timeout(url):
return url.startswith(("redis://", "rediss://", "sqs://", "sentinel://"))
@@ -79,9 +91,12 @@ def get_default_celery_config(team_conf) -> dict[str, Any]:
else:
log.debug("Value for celery result_backend not found. Using
sql_alchemy_conn with db+ prefix.")
sql_alchemy_conn = team_conf.get("database", "SQL_ALCHEMY_CONN")
- # In SQLAlchemy 2.1 the default PostgreSQL driver changed from
psycopg2 to psycopg (v3).
- # To maintain existing behavior, we explicitly specify psycopg2 for
driverless PostgreSQL URLs.
- sql_alchemy_conn = sql_alchemy_conn.replace("postgresql://",
"postgresql+psycopg2://", 1)
+ # Airflow's sync metadata-database driver default is psycopg (v3);
mirror that explicitly
+ # for driverless PostgreSQL URLs instead of relying on SQLAlchemy's
own default. Fall back to
+ # psycopg2 when psycopg/SQLAlchemy 2.0 aren't both available, matching
PostgresHook's own
+ # USE_PSYCOPG3 gating so this doesn't break environments still on the
older combination.
+ target_scheme = "postgresql+psycopg://" if _USE_PSYCOPG3 else
"postgresql+psycopg2://"
+ sql_alchemy_conn = sql_alchemy_conn.replace("postgresql://",
target_scheme, 1)
result_backend = f"db+{sql_alchemy_conn}"
# Handle result backend transport options (for Redis Sentinel support)
diff --git a/providers/celery/src/airflow/providers/celery/get_provider_info.py
b/providers/celery/src/airflow/providers/celery/get_provider_info.py
index a715038d841..4a79cc58821 100644
--- a/providers/celery/src/airflow/providers/celery/get_provider_info.py
+++ b/providers/celery/src/airflow/providers/celery/get_provider_info.py
@@ -130,7 +130,7 @@ def get_provider_info():
"version_added": None,
"type": "string",
"sensitive": True,
- "example":
"db+postgresql+psycopg2://postgres:airflow@postgres/airflow",
+ "example":
"db+postgresql+psycopg://postgres:airflow@postgres/airflow",
"default": None,
},
"result_backend_sqlalchemy_engine_options": {
diff --git
a/providers/celery/tests/unit/celery/executors/test_celery_executor.py
b/providers/celery/tests/unit/celery/executors/test_celery_executor.py
index f66c8e81ab3..28b08405902 100644
--- a/providers/celery/tests/unit/celery/executors/test_celery_executor.py
+++ b/providers/celery/tests/unit/celery/executors/test_celery_executor.py
@@ -847,6 +847,32 @@ def
test_result_backend_transport_options_with_multiple_options():
assert result_backend_opts["master_name"] == "mymaster"
+@conf_vars(
+ {
+ ("celery", "result_backend"): None,
+ ("database", "sql_alchemy_conn"): "postgresql://user:pass@host/db",
+ }
+)
+def
test_result_backend_derived_from_sql_alchemy_conn_uses_psycopg(monkeypatch):
+ """A driverless sql_alchemy_conn must derive a psycopg (v3)
result_backend, not psycopg2."""
+ monkeypatch.setattr(default_celery, "_USE_PSYCOPG3", True)
+ config = default_celery.get_default_celery_config(conf)
+ assert config["result_backend"] ==
"db+postgresql+psycopg://user:pass@host/db"
+
+
+@conf_vars(
+ {
+ ("celery", "result_backend"): None,
+ ("database", "sql_alchemy_conn"): "postgresql://user:pass@host/db",
+ }
+)
+def test_result_backend_falls_back_to_psycopg2_without_psycopg3(monkeypatch):
+ """Without psycopg/SQLAlchemy 2.0 available, the derivation must fall back
to psycopg2."""
+ monkeypatch.setattr(default_celery, "_USE_PSYCOPG3", False)
+ config = default_celery.get_default_celery_config(conf)
+ assert config["result_backend"] ==
"db+postgresql+psycopg2://user:pass@host/db"
+
+
@conf_vars({("celery_result_backend_transport_options", "sentinel_kwargs"):
"invalid_json"})
def test_result_backend_sentinel_kwargs_invalid_json():
"""Test that invalid JSON in sentinel_kwargs raises an error."""
diff --git a/providers/common/sql/src/airflow/providers/common/sql/hooks/sql.py
b/providers/common/sql/src/airflow/providers/common/sql/hooks/sql.py
index 86fb4a466f9..8b5a4d2ddfd 100644
--- a/providers/common/sql/src/airflow/providers/common/sql/hooks/sql.py
+++ b/providers/common/sql/src/airflow/providers/common/sql/hooks/sql.py
@@ -22,6 +22,7 @@ from collections.abc import Callable, Generator, Iterable,
Mapping, MutableMappi
from contextlib import closing, contextmanager, suppress
from datetime import datetime
from functools import cached_property
+from importlib.util import find_spec
from typing import TYPE_CHECKING, Any, Literal, Protocol, TypeVar, cast,
overload
from urllib.parse import urlparse
@@ -66,6 +67,16 @@ if TYPE_CHECKING:
T = TypeVar("T")
SQL_PLACEHOLDERS = frozenset({"%s", "?"})
+
+
+def _is_sqlalchemy_2() -> bool:
+ """Whether the installed SQLAlchemy has the native psycopg (v3) dialect
(added in 2.0)."""
+ from packaging.version import Version
+ from sqlalchemy import __version__ as sqlalchemy_version
+
+ return Version(sqlalchemy_version) >= Version("2.0.0")
+
+
WARNING_MESSAGE = """Import of {} from the
'airflow.providers.common.sql.hooks' module is deprecated and will
be removed in the future. Please import it from
'airflow.providers.common.sql.hooks.handlers'."""
@@ -334,7 +345,31 @@ class DbApiHook(BaseHook):
self.log.debug("url: %s", url)
self.log.debug("engine_kwargs: %s", engine_kwargs)
- return create_engine(url=url, **engine_kwargs)
+ try:
+ return create_engine(url=url, **engine_kwargs)
+ except (ImportError, ModuleNotFoundError):
+ # SQLAlchemy resolves a bare "postgresql" scheme to the psycopg2
DB-API by default.
+ # Retry with a driver this hook can actually resolve, since
psycopg2 may not be
+ # installed now that it's an optional extra of
apache-airflow-providers-postgres.
+ parsed_url = make_url(url) if make_url is not None else None
+ if parsed_url is None or parsed_url.drivername != "postgresql":
+ raise
+ # psycopg (v3) may be importable where SQLAlchemy is pinned below
2.0, which has no
+ # native "postgresql+psycopg" dialect and would raise
NoSuchModuleError — so only
+ # prefer it on SQLAlchemy 2.0+; otherwise fall back to psycopg2.
+ if find_spec("psycopg") is not None and _is_sqlalchemy_2():
+ self.log.info(
+ "SQLAlchemy could not load a DB-API driver for the bare
'postgresql://' URL; "
+ "retrying with psycopg (v3) ('postgresql+psycopg')."
+ )
+ return
create_engine(url=parsed_url.set(drivername="postgresql+psycopg"),
**engine_kwargs)
+ if find_spec("psycopg2") is not None:
+ self.log.info(
+ "SQLAlchemy could not load a DB-API driver for the bare
'postgresql://' URL; "
+ "retrying with 'postgresql+psycopg2'."
+ )
+ return
create_engine(url=parsed_url.set(drivername="postgresql+psycopg2"),
**engine_kwargs)
+ raise
@cached_property
def inspector(self) -> Inspector:
diff --git a/providers/common/sql/tests/unit/common/sql/hooks/test_sql.py
b/providers/common/sql/tests/unit/common/sql/hooks/test_sql.py
index d88cc3abbc8..c2fbfec3cbb 100644
--- a/providers/common/sql/tests/unit/common/sql/hooks/test_sql.py
+++ b/providers/common/sql/tests/unit/common/sql/hooks/test_sql.py
@@ -334,6 +334,98 @@ class TestDbApiHook:
assert isinstance(df, expected_type)
+class TestDbApiHookGetSqlalchemyEngine:
+ """SQLAlchemy resolves a bare ``postgresql`` scheme to the psycopg2 DB-API
by default, which
+ may not be installed now that it's an optional extra of
apache-airflow-providers-postgres."""
+
+ @pytest.mark.db_test
+ @patch("airflow.providers.common.sql.hooks.sql._is_sqlalchemy_2",
return_value=True)
+ @patch("airflow.providers.common.sql.hooks.sql.find_spec")
+ @patch("airflow.providers.common.sql.hooks.sql.create_engine")
+ def test_falls_back_to_psycopg_when_psycopg2_import_fails(
+ self, mock_create_engine, mock_find_spec, mock_is_sqlalchemy_2, caplog
+ ):
+ mock_find_spec.return_value = object()
+ mock_create_engine.side_effect = [ModuleNotFoundError("No module named
'psycopg2'"), MagicMock()]
+ dbapi_hook = mock_db_hook(DbApiHook, conn_params={"conn_type":
"postgresql"})
+
+ dbapi_hook.get_sqlalchemy_engine()
+
+ assert mock_create_engine.call_count == 2
+ retried_url = mock_create_engine.call_args.kwargs["url"]
+ assert retried_url.drivername == "postgresql+psycopg"
+ assert (
+ "SQLAlchemy could not load a DB-API driver for the bare
'postgresql://' URL; "
+ "retrying with psycopg (v3) ('postgresql+psycopg')."
+ ) in caplog.text
+
+ @pytest.mark.db_test
+ @patch("airflow.providers.common.sql.hooks.sql._is_sqlalchemy_2",
return_value=False)
+ @patch("airflow.providers.common.sql.hooks.sql.find_spec")
+ @patch("airflow.providers.common.sql.hooks.sql.create_engine")
+ def test_falls_back_to_psycopg2_when_sqlalchemy_below_2(
+ self, mock_create_engine, mock_find_spec, mock_is_sqlalchemy_2, caplog
+ ):
+ # SQLAlchemy 1.4 (shipped with Airflow 2.11) has no native
"postgresql+psycopg" dialect,
+ # so even when psycopg (v3) is importable the retry must use psycopg2,
not psycopg.
+ mock_find_spec.return_value = object()
+ mock_create_engine.side_effect = [ModuleNotFoundError("No module named
'psycopg2'"), MagicMock()]
+ dbapi_hook = mock_db_hook(DbApiHook, conn_params={"conn_type":
"postgresql"})
+
+ dbapi_hook.get_sqlalchemy_engine()
+
+ assert mock_create_engine.call_count == 2
+ retried_url = mock_create_engine.call_args.kwargs["url"]
+ assert retried_url.drivername == "postgresql+psycopg2"
+ assert (
+ "SQLAlchemy could not load a DB-API driver for the bare
'postgresql://' URL; "
+ "retrying with 'postgresql+psycopg2'."
+ ) in caplog.text
+
+ @pytest.mark.db_test
+ @patch("airflow.providers.common.sql.hooks.sql.find_spec")
+ @patch("airflow.providers.common.sql.hooks.sql.create_engine")
+ def test_falls_back_to_psycopg2_when_psycopg_unavailable(
+ self, mock_create_engine, mock_find_spec, caplog
+ ):
+ mock_find_spec.side_effect = lambda name: None if name == "psycopg"
else object()
+ mock_create_engine.side_effect = [ModuleNotFoundError("No module named
'psycopg2'"), MagicMock()]
+ dbapi_hook = mock_db_hook(DbApiHook, conn_params={"conn_type":
"postgresql"})
+
+ dbapi_hook.get_sqlalchemy_engine()
+
+ assert mock_create_engine.call_count == 2
+ retried_url = mock_create_engine.call_args.kwargs["url"]
+ assert retried_url.drivername == "postgresql+psycopg2"
+ assert (
+ "SQLAlchemy could not load a DB-API driver for the bare
'postgresql://' URL; "
+ "retrying with 'postgresql+psycopg2'."
+ ) in caplog.text
+
+ @pytest.mark.db_test
+ @patch("airflow.providers.common.sql.hooks.sql.find_spec",
return_value=None)
+ @patch("airflow.providers.common.sql.hooks.sql.create_engine")
+ def test_reraises_when_no_postgres_driver_is_available(self,
mock_create_engine, mock_find_spec):
+ original_error = ModuleNotFoundError("No module named 'psycopg2'")
+ mock_create_engine.side_effect = original_error
+ dbapi_hook = mock_db_hook(DbApiHook, conn_params={"conn_type":
"postgresql"})
+
+ with pytest.raises(ModuleNotFoundError, match="psycopg2"):
+ dbapi_hook.get_sqlalchemy_engine()
+ mock_create_engine.assert_called_once()
+
+ @pytest.mark.db_test
+ @patch("airflow.providers.common.sql.hooks.sql.create_engine")
+ def test_does_not_retry_for_non_postgres_schemes(self, mock_create_engine):
+ original_error = ModuleNotFoundError("No module named
'some_other_driver'")
+ mock_create_engine.side_effect = original_error
+ dbapi_hook = mock_db_hook(DbApiHook, conn_params={"conn_type":
"mysql"})
+
+ with pytest.raises(ModuleNotFoundError, match="some_other_driver"):
+ dbapi_hook.get_sqlalchemy_engine()
+ mock_create_engine.assert_called_once()
+
+
class TestDbApiHookGetTableSchema:
@pytest.mark.db_test
def test_get_table_schema(self):
diff --git
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_postgres.py
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_postgres.py
index 76ae2f277bc..61e08326cf7 100644
---
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_postgres.py
+++
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_postgres.py
@@ -22,13 +22,10 @@ from __future__ import annotations
from functools import cached_property
from typing import TYPE_CHECKING
-from psycopg2.extensions import register_adapter
-from psycopg2.extras import Json
-
from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook
from airflow.providers.google.cloud.transfers.bigquery_to_sql import
BigQueryToSqlBaseOperator
from airflow.providers.google.cloud.utils.bigquery_get_data import
bigquery_get_data
-from airflow.providers.postgres.hooks.postgres import PostgresHook
+from airflow.providers.postgres.hooks.postgres import USE_PSYCOPG3,
PostgresHook
if TYPE_CHECKING:
from airflow.providers.common.compat.sdk import Context
@@ -81,8 +78,12 @@ class BigQueryToPostgresOperator(BigQueryToSqlBaseOperator):
@cached_property
def postgres_hook(self) -> PostgresHook:
- register_adapter(list, Json)
- register_adapter(dict, Json)
+ if not USE_PSYCOPG3:
+ from psycopg2.extensions import register_adapter
+ from psycopg2.extras import Json
+
+ register_adapter(list, Json)
+ register_adapter(dict, Json)
return PostgresHook(database=self.database,
postgres_conn_id=self.postgres_conn_id)
def get_sql_hook(self) -> PostgresHook:
diff --git
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
index 85b301fc59e..b44cefb1f30 100644
---
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
+++
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
@@ -17,6 +17,8 @@
# under the License.
from __future__ import annotations
+import importlib
+import sys
from unittest import mock
from unittest.mock import MagicMock
@@ -148,8 +150,13 @@ class TestBigQueryToPostgresOperator:
@mock.patch(
"airflow.providers.google.common.hooks.base_google.GoogleBaseHook.get_credentials_and_project_id"
)
-
@mock.patch("airflow.providers.google.cloud.transfers.bigquery_to_postgres.register_adapter")
- def test_adapters_to_json_registered(self, mock_register_adapter,
mock_get_creds, mock_get_client):
+ @mock.patch("psycopg2.extensions.register_adapter")
+ def test_adapters_to_json_registered_on_psycopg2_path(
+ self, mock_register_adapter, mock_get_creds, mock_get_client,
monkeypatch
+ ):
+ monkeypatch.setattr(
+
"airflow.providers.google.cloud.transfers.bigquery_to_postgres.USE_PSYCOPG3",
False
+ )
mock_get_creds.return_value = (None, TEST_PROJECT)
client = MagicMock()
client.list_rows.return_value = []
@@ -166,6 +173,32 @@ class TestBigQueryToPostgresOperator:
mock_register_adapter.assert_any_call(list, Json)
mock_register_adapter.assert_any_call(dict, Json)
+
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.BigQueryHook.get_client")
+ @mock.patch(
+
"airflow.providers.google.common.hooks.base_google.GoogleBaseHook.get_credentials_and_project_id"
+ )
+ @mock.patch("psycopg2.extensions.register_adapter")
+ def test_adapters_not_registered_on_psycopg3_path(
+ self, mock_register_adapter, mock_get_creds, mock_get_client,
monkeypatch
+ ):
+ monkeypatch.setattr(
+
"airflow.providers.google.cloud.transfers.bigquery_to_postgres.USE_PSYCOPG3",
True
+ )
+ mock_get_creds.return_value = (None, TEST_PROJECT)
+ client = MagicMock()
+ client.list_rows.return_value = []
+ mock_get_client.return_value = client
+
+ operator = BigQueryToPostgresOperator(
+ task_id=TASK_ID,
+ dataset_table=f"{TEST_DATASET}.{TEST_TABLE_ID}",
+ target_table_name=TEST_DESTINATION_TABLE,
+ replace=False,
+ )
+ operator.postgres_hook
+
+ mock_register_adapter.assert_not_called()
+
@mock.patch("airflow.providers.google.cloud.transfers.bigquery_to_postgres.PostgresHook")
@mock.patch("airflow.providers.google.cloud.transfers.bigquery_to_sql.BigQueryHook")
def test_get_openlineage_facets_on_complete_no_selected_fields(self,
mock_bq_hook, mock_postgres_hook):
@@ -254,3 +287,17 @@ class TestBigQueryToPostgresOperator:
assert "columnLineage" in output_ds.facets
col_lineage = output_ds.facets["columnLineage"]
assert set(col_lineage.fields.keys()) == {"id", "name"}
+
+
+def test_bigquery_to_postgres_module_imports_without_psycopg2(monkeypatch):
+ """The module must import cleanly even when psycopg2 isn't installed."""
+ monkeypatch.setitem(sys.modules, "psycopg2", None)
+ monkeypatch.setitem(sys.modules, "psycopg2.extensions", None)
+ monkeypatch.setitem(sys.modules, "psycopg2.extras", None)
+ module_name =
"airflow.providers.google.cloud.transfers.bigquery_to_postgres"
+ monkeypatch.delitem(sys.modules, module_name, raising=False)
+ try:
+ module = importlib.import_module(module_name)
+ assert module.BigQueryToPostgresOperator is not None
+ finally:
+ monkeypatch.delitem(sys.modules, module_name, raising=False)
diff --git a/providers/pgvector/pyproject.toml
b/providers/pgvector/pyproject.toml
index 4bd7a8a766c..8cfd22cf8dc 100644
--- a/providers/pgvector/pyproject.toml
+++ b/providers/pgvector/pyproject.toml
@@ -77,6 +77,7 @@ dev = [
"apache-airflow",
"apache-airflow-task-sdk",
"apache-airflow-devel-common",
+ "apache-airflow-providers-common-compat",
"apache-airflow-providers-common-sql",
"apache-airflow-providers-postgres",
# Additional devel dependencies (do not remove this line and add extra
development dependencies)
diff --git
a/providers/pgvector/src/airflow/providers/pgvector/operators/pgvector.py
b/providers/pgvector/src/airflow/providers/pgvector/operators/pgvector.py
index ae48a377ea7..b45fdb65f89 100644
--- a/providers/pgvector/src/airflow/providers/pgvector/operators/pgvector.py
+++ b/providers/pgvector/src/airflow/providers/pgvector/operators/pgvector.py
@@ -17,8 +17,7 @@
# under the License.
from __future__ import annotations
-from pgvector.psycopg2 import register_vector
-
+from airflow.providers.common.compat.sdk import
AirflowOptionalProviderFeatureException
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
@@ -41,6 +40,24 @@ class PgVectorIngestOperator(SQLExecuteQueryOperator):
def _register_vector(self) -> None:
"""Register the vector type with your connection."""
+ from airflow.providers.postgres.hooks.postgres import USE_PSYCOPG3
+
+ if USE_PSYCOPG3:
+ try:
+ from pgvector.psycopg import register_vector
+ except (ImportError, ModuleNotFoundError) as err:
+ raise AirflowOptionalProviderFeatureException(
+ "pgvector's psycopg (v3) integration is not installed.
Please install it with "
+ "`pip install --upgrade pgvector`."
+ ) from err
+ else:
+ try:
+ from pgvector.psycopg2 import register_vector
+ except (ImportError, ModuleNotFoundError) as err:
+ raise AirflowOptionalProviderFeatureException(
+ "psycopg2 is not installed. Please install it with "
+ "`pip install
apache-airflow-providers-postgres[psycopg2]`."
+ ) from err
conn = self.get_db_hook().get_conn()
register_vector(conn)
diff --git a/providers/pgvector/tests/unit/pgvector/operators/test_pgvector.py
b/providers/pgvector/tests/unit/pgvector/operators/test_pgvector.py
index 6fe144d5bee..0dec5bd5826 100644
--- a/providers/pgvector/tests/unit/pgvector/operators/test_pgvector.py
+++ b/providers/pgvector/tests/unit/pgvector/operators/test_pgvector.py
@@ -16,10 +16,13 @@
# under the License.
from __future__ import annotations
+import importlib
+import sys
from unittest.mock import Mock, patch
import pytest
+from airflow.providers.common.compat.sdk import
AirflowOptionalProviderFeatureException
from airflow.providers.pgvector.operators.pgvector import
PgVectorIngestOperator
@@ -32,10 +35,24 @@ def pg_vector_ingest_operator():
)
-@patch("airflow.providers.pgvector.operators.pgvector.register_vector")
+@patch("airflow.providers.postgres.hooks.postgres.USE_PSYCOPG3", False)
@patch("airflow.providers.pgvector.operators.pgvector.PgVectorIngestOperator.get_db_hook")
-def test_register_vector(mock_get_db_hook, mock_register_vector,
pg_vector_ingest_operator):
- # Create a mock database connection
+def test_register_vector_psycopg2(mock_get_db_hook, monkeypatch,
pg_vector_ingest_operator):
+ # psycopg2 is optional (see providers/postgres's [psycopg2] extra), so
patch the module
+ # via sys.modules instead of `@patch("pgvector.psycopg2...")`, which would
import it for real.
+ mock_psycopg2 = Mock()
+ monkeypatch.setitem(sys.modules, "pgvector.psycopg2", mock_psycopg2)
+ mock_db_hook = Mock()
+ mock_get_db_hook.return_value = mock_db_hook
+
+ pg_vector_ingest_operator._register_vector()
+ mock_psycopg2.register_vector.assert_called_with(mock_db_hook.get_conn())
+
+
+@patch("airflow.providers.postgres.hooks.postgres.USE_PSYCOPG3", True)
+@patch("pgvector.psycopg.register_vector")
+@patch("airflow.providers.pgvector.operators.pgvector.PgVectorIngestOperator.get_db_hook")
+def test_register_vector_psycopg3(mock_get_db_hook, mock_register_vector,
pg_vector_ingest_operator):
mock_db_hook = Mock()
mock_get_db_hook.return_value = mock_db_hook
@@ -43,14 +60,44 @@ def test_register_vector(mock_get_db_hook,
mock_register_vector, pg_vector_inges
mock_register_vector.assert_called_with(mock_db_hook.get_conn())
-@patch("airflow.providers.pgvector.operators.pgvector.register_vector")
+@patch("airflow.providers.postgres.hooks.postgres.USE_PSYCOPG3", False)
@patch("airflow.providers.pgvector.operators.pgvector.SQLExecuteQueryOperator.execute")
@patch("airflow.providers.pgvector.operators.pgvector.PgVectorIngestOperator.get_db_hook")
def test_execute(
- mock_get_db_hook, mock_execute_query_operator_execute,
mock_register_vector, pg_vector_ingest_operator
+ mock_get_db_hook, mock_execute_query_operator_execute, monkeypatch,
pg_vector_ingest_operator
):
+ # psycopg2 is optional (see providers/postgres's [psycopg2] extra), so
patch the module
+ # via sys.modules instead of `@patch("pgvector.psycopg2...")`, which would
import it for real.
+ monkeypatch.setitem(sys.modules, "pgvector.psycopg2", Mock())
mock_db_hook = Mock()
mock_get_db_hook.return_value = mock_db_hook
pg_vector_ingest_operator.execute(None)
mock_execute_query_operator_execute.assert_called_once()
+
+
+@patch("airflow.providers.postgres.hooks.postgres.USE_PSYCOPG3", False)
+def test_register_vector_raises_clear_error_without_psycopg2(monkeypatch,
pg_vector_ingest_operator):
+ monkeypatch.setitem(sys.modules, "pgvector.psycopg2", None)
+ with pytest.raises(AirflowOptionalProviderFeatureException,
match="psycopg2 is not installed"):
+ pg_vector_ingest_operator._register_vector()
+
+
+@patch("airflow.providers.postgres.hooks.postgres.USE_PSYCOPG3", True)
+def test_register_vector_raises_clear_error_without_psycopg3(monkeypatch,
pg_vector_ingest_operator):
+ monkeypatch.setitem(sys.modules, "pgvector.psycopg", None)
+ with pytest.raises(AirflowOptionalProviderFeatureException,
match=r"psycopg \(v3\) integration"):
+ pg_vector_ingest_operator._register_vector()
+
+
+def test_pgvector_module_imports_without_psycopg2(monkeypatch):
+ """The module must import cleanly even when psycopg2/pgvector.psycopg2
isn't installed."""
+ monkeypatch.setitem(sys.modules, "psycopg2", None)
+ monkeypatch.setitem(sys.modules, "pgvector.psycopg2", None)
+ module_name = "airflow.providers.pgvector.operators.pgvector"
+ monkeypatch.delitem(sys.modules, module_name, raising=False)
+ try:
+ module = importlib.import_module(module_name)
+ assert module.PgVectorIngestOperator is not None
+ finally:
+ monkeypatch.delitem(sys.modules, module_name, raising=False)
diff --git a/providers/postgres/README.rst b/providers/postgres/README.rst
index b499cc4b8ca..15535d883fd 100644
--- a/providers/postgres/README.rst
+++ b/providers/postgres/README.rst
@@ -56,10 +56,10 @@ PIP package Version required
``apache-airflow`` ``>=2.11.0``
``apache-airflow-providers-common-compat`` ``>=1.12.0``
``apache-airflow-providers-common-sql`` ``>=1.32.0``
-``psycopg2-binary`` ``>=2.9.9; python_version <
"3.13"``
-``psycopg2-binary`` ``>=2.9.10; python_version >=
"3.13"``
``psycopg[binary]`` ``>=3.2.9; python_version <
"3.14"``
``psycopg[binary]`` ``>=3.3.3; python_version >=
"3.14"``
+``psycopg2-binary`` ``>=2.9.9; python_version <
"3.13"``
+``psycopg2-binary`` ``>=2.9.10; python_version >=
"3.13"``
``asyncpg`` ``>=0.30.0``
==========================================
======================================
@@ -96,6 +96,7 @@ Extra Dependencies
``openlineage`` ``apache-airflow-providers-openlineage``
``pandas`` ``pandas>=2.1.2; python_version <"3.13"``,
``pandas>=2.2.3; python_version >="3.13" and python_version <"3.14"``,
``pandas>=2.3.3; python_version >="3.14"``
``polars`` ``polars>=1.26.0``
+``psycopg2`` ``psycopg2-binary>=2.9.9; python_version < "3.13"``,
``psycopg2-binary>=2.9.10; python_version >= "3.13"``
``sqlalchemy`` ``sqlalchemy>=1.4.54``
===================
============================================================================================================================================================
diff --git a/providers/postgres/docs/changelog.rst
b/providers/postgres/docs/changelog.rst
index 7198d5756d4..48518f1adc3 100644
--- a/providers/postgres/docs/changelog.rst
+++ b/providers/postgres/docs/changelog.rst
@@ -49,6 +49,20 @@ Breaking changes
To keep using asyncpg on Airflow 3.4.0+, set
``[database] sql_alchemy_conn_async = postgresql+asyncpg://...``
explicitly.
+.. note::
+ On Airflow 3.4.0 and later, ``psycopg`` (psycopg3) becomes the default
synchronous Postgres
+ driver, mirroring the async default above. ``psycopg[binary]`` is
installed by default to serve it.
+
+ ``psycopg2-binary`` remains installed by default as well. This provider
still supports Airflow
+ cores older than 3.4.0, which normalize the synchronous metadata-database
connection to
+ ``postgresql+psycopg2://`` and have no psycopg fallback; Airflow 2.11
additionally ships
+ SQLAlchemy 1.4, which has no native ``psycopg`` (v3) dialect at all — so
psycopg2 must stay
+ present for those cores to work. It will become opt-in-only (via the
existing ``[psycopg2]``
+ extra) in a future release once the minimum supported Airflow is 3.4.0.
+
+ To keep using psycopg2 on Airflow 3.4.0+, set
+ ``[database] sql_alchemy_conn = postgresql+psycopg2://...`` explicitly.
+
* ``Switch the default async Postgres driver from asyncpg to psycopg3
(#69089)``
.. Below changes are excluded from the changelog. Move them to
diff --git a/providers/postgres/docs/index.rst
b/providers/postgres/docs/index.rst
index 71ed92201ca..a236c0a4afc 100644
--- a/providers/postgres/docs/index.rst
+++ b/providers/postgres/docs/index.rst
@@ -104,10 +104,10 @@ PIP package Version
required
``apache-airflow`` ``>=2.11.0``
``apache-airflow-providers-common-compat`` ``>=1.12.0``
``apache-airflow-providers-common-sql`` ``>=1.32.0``
-``psycopg2-binary`` ``>=2.9.9; python_version <
"3.13"``
-``psycopg2-binary`` ``>=2.9.10; python_version >=
"3.13"``
``psycopg[binary]`` ``>=3.2.9; python_version <
"3.14"``
``psycopg[binary]`` ``>=3.3.3; python_version >=
"3.14"``
+``psycopg2-binary`` ``>=2.9.9; python_version <
"3.13"``
+``psycopg2-binary`` ``>=2.9.10; python_version >=
"3.13"``
``asyncpg`` ``>=0.30.0``
==========================================
======================================
@@ -152,6 +152,7 @@ Extra Dependencies
``openlineage`` ``apache-airflow-providers-openlineage``
``pandas`` ``pandas>=2.1.2; python_version <"3.13"``,
``pandas>=2.2.3; python_version >="3.13" and python_version <"3.14"``,
``pandas>=2.3.3; python_version >="3.14"``
``polars`` ``polars>=1.26.0``
+``psycopg2`` ``psycopg2-binary>=2.9.9; python_version < '3.13'``,
``psycopg2-binary>=2.9.10; python_version >= '3.13'``
``sqlalchemy`` ``sqlalchemy>=1.4.54``
===================
============================================================================================================================================================
diff --git a/providers/postgres/pyproject.toml
b/providers/postgres/pyproject.toml
index 71a9a295b85..1898cd5dc40 100644
--- a/providers/postgres/pyproject.toml
+++ b/providers/postgres/pyproject.toml
@@ -62,18 +62,18 @@ dependencies = [
"apache-airflow>=2.11.0",
"apache-airflow-providers-common-compat>=1.12.0",
"apache-airflow-providers-common-sql>=1.32.0",
- # psycopg2 remains the sync driver for now; its removal and the sync
migration to
- # psycopg3 are tracked at https://github.com/apache/airflow/issues/68453
- "psycopg2-binary>=2.9.9; python_version < '3.13'",
- "psycopg2-binary>=2.9.10; python_version >= '3.13'",
"psycopg[binary]>=3.2.9; python_version < '3.14'",
"psycopg[binary]>=3.3.3; python_version >= '3.14'",
- # asyncpg must stay a default dependency while this provider still
supports Airflow
- # cores older than 3.4.0: those cores derive the async metadata-DB URL as
- # ``postgresql+asyncpg://`` unconditionally and have no psycopg_async
fallback, so
- # dropping asyncpg breaks their async engine. Make it opt-in-only (via the
``[asyncpg]``
- # extra) only once the minimum supported Airflow is >= 3.4.0; tracked at
- # https://github.com/apache/airflow/issues/68453
+ # psycopg2-binary and asyncpg must stay default dependencies while this
provider still
+ # supports Airflow cores older than 3.4.0. Those cores normalize the sync
metadata-DB conn
+ # to ``postgresql+psycopg2://`` (see core
``_upgrade_postgres_metastore_conn``) and derive
+ # the async URL as ``postgresql+asyncpg://`` unconditionally, with no
psycopg fallback;
+ # Airflow 2.11 additionally ships SQLAlchemy 1.4, which has no native
``psycopg`` (v3)
+ # dialect at all. The psycopg3 defaults only take effect on 3.4.0+, so
both drivers become
+ # opt-in-only (via the ``[psycopg2]`` / ``[asyncpg]`` extras) together
once the minimum
+ # supported Airflow is >= 3.4.0; tracked at
https://github.com/apache/airflow/issues/68453
+ "psycopg2-binary>=2.9.9; python_version < '3.13'",
+ "psycopg2-binary>=2.9.10; python_version >= '3.13'",
"asyncpg>=0.30.0",
]
@@ -100,6 +100,10 @@ dependencies = [
"polars" = [
"polars>=1.26.0"
]
+"psycopg2" = [
+ "psycopg2-binary>=2.9.9; python_version < '3.13'",
+ "psycopg2-binary>=2.9.10; python_version >= '3.13'",
+]
"sqlalchemy" = [
"sqlalchemy>=1.4.54"
]
diff --git
a/providers/postgres/src/airflow/providers/postgres/hooks/postgres.py
b/providers/postgres/src/airflow/providers/postgres/hooks/postgres.py
index be6519ddade..3821779e261 100644
--- a/providers/postgres/src/airflow/providers/postgres/hooks/postgres.py
+++ b/providers/postgres/src/airflow/providers/postgres/hooks/postgres.py
@@ -18,14 +18,12 @@
from __future__ import annotations
import os
-from collections.abc import Iterable, Mapping
+from collections.abc import Callable, Iterable, Mapping
from contextlib import closing
from copy import deepcopy
-from typing import TYPE_CHECKING, Any, Literal, Protocol, TypeAlias, cast,
overload
+from typing import TYPE_CHECKING, Any, Literal, NoReturn, Protocol, TypeAlias,
cast, overload
from more_itertools import chunked
-from psycopg2 import connect as ppg2_connect
-from psycopg2.extras import DictCursor, NamedTupleCursor, RealDictCursor,
execute_values
from airflow.providers.common.compat.sdk import (
AirflowException,
@@ -54,9 +52,35 @@ if USE_PSYCOPG3:
from psycopg.rows import dict_row, namedtuple_row
from psycopg.types.json import register_default_adapters
+try:
+ import psycopg2 as _psycopg2
+ import psycopg2.extras as _psycopg2_extras
+except (ImportError, ModuleNotFoundError):
+ _psycopg2 = None
+ _psycopg2_extras = None
+
+ppg2_connect: Callable[..., Any] | None = _psycopg2.connect if _psycopg2 else
None
+DictCursor: type | None = _psycopg2_extras.DictCursor if _psycopg2_extras else
None
+NamedTupleCursor: type | None = _psycopg2_extras.NamedTupleCursor if
_psycopg2_extras else None
+RealDictCursor: type | None = _psycopg2_extras.RealDictCursor if
_psycopg2_extras else None
+execute_values: Callable[..., Any] | None = _psycopg2_extras.execute_values if
_psycopg2_extras else None
+
+
+def _require_psycopg2() -> NoReturn:
+ raise AirflowOptionalProviderFeatureException(
+ "psycopg2 is not installed. Please install it with "
+ "`pip install apache-airflow-providers-postgres[psycopg2]`."
+ )
+
+
if TYPE_CHECKING:
from pandas import DataFrame as PandasDataFrame
from polars import DataFrame as PolarsDataFrame
+ from psycopg2.extras import (
+ DictCursor as _DictCursorType,
+ NamedTupleCursor as _NamedTupleCursorType,
+ RealDictCursor as _RealDictCursorType,
+ )
from sqlalchemy.engine import URL
from airflow.providers.common.sql.dialects.dialect import Dialect
@@ -65,7 +89,7 @@ if TYPE_CHECKING:
if USE_PSYCOPG3:
from psycopg.errors import Diagnostic
- CursorType: TypeAlias = DictCursor | RealDictCursor | NamedTupleCursor
+ CursorType: TypeAlias = _DictCursorType | _RealDictCursorType |
_NamedTupleCursorType
CursorRow: TypeAlias = dict[str, Any] | tuple[Any, ...]
@@ -204,6 +228,9 @@ class PostgresHook(DbApiHook):
valid_cursors = "dictcursor, namedtuplecursor"
raise ValueError(f"Invalid cursor passed {_cursor}. Valid options
are: {valid_cursors}")
+ if DictCursor is None:
+ _require_psycopg2()
+
cursor_types = {
"dictcursor": DictCursor,
"realdictcursor": RealDictCursor,
@@ -235,6 +262,9 @@ class PostgresHook(DbApiHook):
return connection
+ if ppg2_connect is None:
+ _require_psycopg2()
+
return ppg2_connect(**conn_args)
def _generate_cursor_name(self):
@@ -691,6 +721,8 @@ class PostgresHook(DbApiHook):
)
# if fast_executemany is enabled with psycopg2, use optimized
execute_values from psycopg
+ if execute_values is None:
+ _require_psycopg2()
self._insert_statement_format = "INSERT INTO {} {} VALUES %s"
nb_rows = 0
diff --git a/providers/postgres/tests/unit/postgres/hooks/test_postgres.py
b/providers/postgres/tests/unit/postgres/hooks/test_postgres.py
index e8c1dd5eff9..a9eff1bcb6f 100644
--- a/providers/postgres/tests/unit/postgres/hooks/test_postgres.py
+++ b/providers/postgres/tests/unit/postgres/hooks/test_postgres.py
@@ -58,6 +58,27 @@ else:
import psycopg2.extras
+def test_hooks_postgres_raises_clear_error_without_psycopg2(monkeypatch):
+ """PostgresHook must fail loudly (not with a bare ImportError or a real
connection attempt)
+ when psycopg2-specific functionality is used but psycopg2 isn't
installed."""
+ import airflow.providers.postgres.hooks.postgres as postgres_module
+
+ monkeypatch.setattr(postgres_module, "USE_PSYCOPG3", False)
+ monkeypatch.setattr(postgres_module, "ppg2_connect", None)
+ monkeypatch.setattr(postgres_module, "DictCursor", None)
+ monkeypatch.setattr(postgres_module, "RealDictCursor", None)
+ monkeypatch.setattr(postgres_module, "NamedTupleCursor", None)
+ monkeypatch.setattr(postgres_module, "execute_values", None)
+
+ hook = postgres_module.PostgresHook.__new__(postgres_module.PostgresHook)
+ with pytest.raises(AirflowOptionalProviderFeatureException,
match="psycopg2 is not installed"):
+ hook._get_cursor("dictcursor")
+ with pytest.raises(AirflowOptionalProviderFeatureException,
match="psycopg2 is not installed"):
+ hook._create_connection({})
+ with pytest.raises(AirflowOptionalProviderFeatureException,
match="psycopg2 is not installed"):
+ hook.insert_rows(table="t", rows=[(1,)], fast_executemany=True)
+
+
@pytest.fixture
def mock_connect(mocker):
"""Mock the connection object according to the correct psycopg version."""
diff --git
a/providers/postgres/tests/unit/postgres/test_provider_dependencies.py
b/providers/postgres/tests/unit/postgres/test_provider_dependencies.py
new file mode 100644
index 00000000000..da83a3a5aeb
--- /dev/null
+++ b/providers/postgres/tests/unit/postgres/test_provider_dependencies.py
@@ -0,0 +1,46 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+from importlib import metadata
+
+import pytest
+
+DISTRIBUTION = "apache-airflow-providers-postgres"
+
+
+def _default_requirements() -> list[str]:
+ """Requirement strings the provider installs by default (i.e. not gated
behind an extra)."""
+ return [req for req in (metadata.requires(DISTRIBUTION) or []) if "extra
==" not in req]
+
+
[email protected]("package", ["psycopg2-binary", "asyncpg"])
+def test_backcompat_driver_stays_a_default_dependency(package):
+ """
+ ``psycopg2-binary`` (sync) and ``asyncpg`` (async) must remain default
dependencies, not
+ optional extras, while the provider supports Airflow cores older than
3.4.0. Those cores
+ normalize the sync metadata-DB conn to ``postgresql+psycopg2://`` and
derive the async URL as
+ ``postgresql+asyncpg://`` with no psycopg fallback (and Airflow 2.11 ships
SQLAlchemy 1.4,
+ which has no ``psycopg`` v3 dialect), so demoting either driver to an
extra silently breaks
+ them. Removal is gated on bumping the floor to 3.4.0; see
+ https://github.com/apache/airflow/issues/68453
+ """
+ defaults = _default_requirements()
+ assert any(req.startswith(package) for req in defaults), (
+ f"{package} must stay a default dependency of {DISTRIBUTION} until the
minimum supported "
+ f"Airflow is >= 3.4.0. Default dependencies were: {defaults}"
+ )
diff --git a/uv.lock b/uv.lock
index 715b7175aaa..045af593989 100644
--- a/uv.lock
+++ b/uv.lock
@@ -7104,6 +7104,7 @@ common-sql = [
dev = [
{ name = "apache-airflow" },
{ name = "apache-airflow-devel-common" },
+ { name = "apache-airflow-providers-common-compat" },
{ name = "apache-airflow-providers-common-sql" },
{ name = "apache-airflow-providers-postgres" },
{ name = "apache-airflow-task-sdk" },
@@ -7126,6 +7127,7 @@ provides-extras = ["common-sql"]
dev = [
{ name = "apache-airflow", editable = "." },
{ name = "apache-airflow-devel-common", editable = "devel-common" },
+ { name = "apache-airflow-providers-common-compat", editable =
"providers/common/compat" },
{ name = "apache-airflow-providers-common-sql", editable =
"providers/common/sql" },
{ name = "apache-airflow-providers-postgres", editable =
"providers/postgres" },
{ name = "apache-airflow-task-sdk", editable = "task-sdk" },
@@ -7201,6 +7203,9 @@ pandas = [
polars = [
{ name = "polars" },
]
+psycopg2 = [
+ { name = "psycopg2-binary" },
+]
sqlalchemy = [
{ name = "sqlalchemy" },
]
@@ -7239,9 +7244,11 @@ requires-dist = [
{ name = "psycopg", extras = ["binary"], marker = "python_full_version >=
'3.14'", specifier = ">=3.3.3" },
{ name = "psycopg2-binary", marker = "python_full_version < '3.13'",
specifier = ">=2.9.9" },
{ name = "psycopg2-binary", marker = "python_full_version >= '3.13'",
specifier = ">=2.9.10" },
+ { name = "psycopg2-binary", marker = "python_full_version >= '3.13' and
extra == 'psycopg2'", specifier = ">=2.9.10" },
+ { name = "psycopg2-binary", marker = "python_full_version < '3.13' and
extra == 'psycopg2'", specifier = ">=2.9.9" },
{ name = "sqlalchemy", marker = "extra == 'sqlalchemy'", specifier =
">=1.4.54" },
]
-provides-extras = ["amazon", "asyncpg", "microsoft-azure", "openlineage",
"pandas", "polars", "sqlalchemy"]
+provides-extras = ["amazon", "asyncpg", "microsoft-azure", "openlineage",
"pandas", "polars", "psycopg2", "sqlalchemy"]
[package.metadata.requires-dev]
dev = [