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

Reply via email to