anmolxlight commented on code in PR #71074:
URL: https://github.com/apache/airflow/pull/71074#discussion_r4152261973


##########
airflow-core/src/airflow/assets/manager.py:
##########
@@ -754,27 +698,124 @@ def _get_or_create_apdr(
         target_partition_date: datetime | None,
         target_dag: DagModel,
         rollup_fingerprint: dict,
-        asset_id: int,
         session: Session,
     ) -> AssetPartitionDagRun:
         """
         Get or create an APDR.
 
-        If 2 processes invoke this method at the same time using the same 
(target_key, target_dag) pair,
-        they may both check the database and, finding no existing APDR, create 
separate instances.
-        This leads to the unintended outcome of having two APDRs created 
instead of one.
-        To resolve this, we add a mutex lock to AssetModel for PostgreSQL and 
MySQL and use
-        AssetPartitionDagRunMutexLock table for SQLite.
+        If 2 processes invoke this method at the same time using the same 
(target_key, target_dag)
+        pair, they may both check the database and, finding no existing APDR, 
attempt to create
+        separate instances. Rather than serializing this find-or-create behind 
a lock, a unique
+        constraint on (target_dag_id, pending_partition_key) — see the 
``AssetPartitionDagRun``
+        docstring — makes the database itself reject the loser's INSERT. The 
loser catches that
+        ``IntegrityError`` and re-selects, working on the winning row instead 
of raising, per the
+        model's "always work on the latest matching APDR record" contract. 
Optimistic and
+        lock-free, this scales with concurrent producer assets without 
contending on a Dag or
+        Asset row that has nothing to do with the (target_key, target_dag) 
pair being deduplicated.
 
         ``rollup_fingerprint`` is the serialized mapper / window definition 
for all partitioned
         assets in the timetable at creation time; the scheduler discards APDRs 
whose stamp no
         longer matches the current timetable's fingerprint (mapper / window 
may have changed).
+        """
+        latest_apdr = cls._get_latest_pending_apdr(
+            target_key=target_key, target_dag_id=target_dag.dag_id, 
session=session
+        )
+        if latest_apdr is not None:
+            cls._reconcile_partition_date(
+                apdr=latest_apdr,
+                target_partition_date=target_partition_date,
+                target_dag_id=target_dag.dag_id,
+                target_key=target_key,
+                session=session,
+            )
+            cls.logger().debug(
+                "Existing APDR found for key %s dag_id %s",
+                target_key,
+                target_dag.dag_id,
+                exc_info=True,
+            )
+            return latest_apdr
+
+        apdr = AssetPartitionDagRun(
+            target_dag_id=target_dag.dag_id,
+            created_dag_run_id=None,
+            partition_key=target_key,
+            pending_partition_key=target_key,
+            partition_date=target_partition_date,
+            rollup_fingerprint=rollup_fingerprint,
+        )
+        try:
+            # A SAVEPOINT scopes the potential IntegrityError so only this 
INSERT is rolled
+            # back on conflict; the caller's surrounding transaction (with any 
other work
+            # already flushed in this scheduler tick) stays intact.
+            with session.begin_nested():
+                session.add(apdr)
+                session.flush()

Review Comment:
   Re-verified against the current head and left the fix in place rather than 
reshuffling the loop.
   
   The concern is about the mechanics of the rollback, and you are right about 
them: `_restore_snapshot()` does `to_expunge = set(self._new) | 
set(session._new)`. What is missing is when the pending rows leave `_new`. 
`SessionTransaction.__init__` calls `_take_snapshot()` for every transaction it 
opens, nested ones included, and for a `BEGIN_NESTED` origin `_take_snapshot()` 
runs `self.session.flush()` before the SAVEPOINT exists. So at the moment 
`_get_or_create_apdr` opens the SAVEPOINT, `session._new` is already empty: 
every Log / PartitionedAssetKeyLog row accumulated earlier in 
`_queue_partitioned_dags` was flushed into the parent transaction and is 
persistent. Only `apdr`, added after the SAVEPOINT opened, is in the nested 
snapshot, so the rollback can only expunge that one.
   
   Pinned by `test_lost_race_does_not_expunge_unflushed_log_rows` 
(`airflow-core/tests/unit/assets/test_manager.py`): an unflushed Log row is 
added before the SAVEPOINT, a colliding INSERT forces the race, and the row is 
asserted to persist afterwards. The comment above `with 
session.begin_nested():` says the same thing for the next reader, and 
`Session.begin_nested()` is what the line does, not an assumption about library 
internals.
   
   Verified on SQLAlchemy 2.0.51 by reading `_take_snapshot` / 
`_restore_snapshot` and by running the repro. Happy to restructure 
`_queue_partitioned_dags` into a per-target savepoint if you would rather not 
rely on that flush; it is not needed for correctness here.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to