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 9e509db2a7b docs - add create_async_metadata_engine example (#71663)
9e509db2a7b is described below
commit 9e509db2a7b4b8141c6f3933bdad5479f280b50f
Author: raphaelauv <[email protected]>
AuthorDate: Wed Sep 9 14:15:19 2026 +0200
docs - add create_async_metadata_engine example (#71663)
Co-authored-by: raphaelauv <[email protected]>
---
.../cluster-policies.rst | 21 +++++++++++++++++++--
1 file changed, 19 insertions(+), 2 deletions(-)
diff --git
a/airflow-core/docs/administration-and-deployment/cluster-policies.rst
b/airflow-core/docs/administration-and-deployment/cluster-policies.rst
index e23fd7ffad2..7bc1ba5a377 100644
--- a/airflow-core/docs/administration-and-deployment/cluster-policies.rst
+++ b/airflow-core/docs/administration-and-deployment/cluster-policies.rst
@@ -210,7 +210,7 @@ Two functions can be overridden:
* ``create_metadata_engine(sql_alchemy_conn, *, engine_args, connect_args) ->
Engine`` — called by
``configure_orm()`` to create the synchronous metadata engine.
-* ``create_async_metadata_engine(sql_alchemy_conn_async, *, connect_args) ->
AsyncEngine`` — called by
+* ``create_async_metadata_engine(sql_alchemy_conn_async, *, connect_args,
engine_args) -> AsyncEngine`` — called by
``_configure_async_session()`` to create the asynchronous metadata engine.
The default implementations call ``sqlalchemy.create_engine`` /
``sqlalchemy.ext.asyncio.create_async_engine``
@@ -231,7 +231,7 @@ Example: registering a ``do_connect`` handler that
refreshes a JWT token before
dbapi_connection.execute(f"SET SESSION AUTHORIZATION '{token}'")
- def create_metadata_engine(sql_alchemy_conn, *, engine_args, connect_args):
+ def create_metadata_engine(sql_alchemy_conn, *, engine_args, connect_args)
-> Engine:
engine = create_engine(
sql_alchemy_conn,
connect_args=connect_args,
@@ -240,3 +240,20 @@ Example: registering a ``do_connect`` handler that
refreshes a JWT token before
)
event.listen(engine, "do_connect", _refresh_jwt)
return engine
+
+
+ def create_async_metadata_engine(sql_alchemy_conn_async, *, connect_args,
engine_args) -> AsyncEngine:
+ connect_args["ssl"] = "require"
+ engine = create_async_engine(
+ sql_alchemy_conn_async,
+ connect_args=connect_args,
+ **engine_args,
+ future=True,
+ )
+
+ @event.listens_for(engine.sync_engine, "do_connect")
+ def provide_token(dialect, conn_rec, cargs, cparams):
+ token = my_token_provider.get_token(user=cparams["user"],
host=cparams["host"], port=cparams["port"])
+ cparams["password"] = token
+
+ return engine