rjgoyln commented on code in PR #73119:
URL: https://github.com/apache/airflow/pull/73119#discussion_r4115150821


##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py:
##########
@@ -124,6 +127,89 @@ def _generate_queued_event_where_clause(
     return where_clause
 
 
+QueuedEventPartitionKeyQuery = Annotated[
+    str | None,
+    Query(
+        description=(
+            "Delete queued events of partitioned assets emitted for this 
partition key whose "
+            "Dag run has not been created yet, instead of queued events of 
non-partitioned assets."
+        ),
+    ),
+]
+
+
+def _delete_pending_partitioned_queued_events(
+    *,
+    partition_key: str,
+    session: Session,
+    asset_id: int | None = None,
+    dag_id: str | None = None,
+    before: datetime | str | None = None,
+    permitted_dag_ids: set[str] | None = None,
+) -> int:
+    """
+    Delete queued partitioned asset events whose Dag run has not been created 
yet.
+
+    ``partition_key`` is matched against the partition key of the asset event, 
as in
+    ``POST /assets/{asset_id}/materialize``, not against the partition key of 
the Dag run
+    the event would create; a partition mapper can make the two differ.
+
+    Only events of an existing ``AssetPartitionDagRun`` whose Dag run has not 
been created
+    yet are deleted; ``PartitionedAssetKeyLog`` rows left behind once their 
Dag run was
+    deleted are history, not queued events. Pending ``AssetPartitionDagRun`` 
rows left
+    without any contributing events are deleted afterwards, so clearing some 
of the events
+    of a partition keeps the partition waiting for them rather than dropping 
the progress
+    made by the others.
+
+    :return: The number of deleted queued events.
+    """
+    pakl_where_clause: list[ColumnElement[bool]] = [
+        PartitionedAssetKeyLog.source_partition_key == partition_key,
+    ]
+    if asset_id is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.asset_id == asset_id)
+    if dag_id is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.target_dag_id == 
dag_id)
+    if before is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.created_at < before)
+    if permitted_dag_ids is not None:
+        
pakl_where_clause.append(PartitionedAssetKeyLog.target_dag_id.in_(permitted_dag_ids))
+    pakl_of_apdr = PartitionedAssetKeyLog.asset_partition_dag_run_id == 
AssetPartitionDagRun.id
+
+    # Nothing is loaded into the session beforehand, so skip synchronizing it; 
the subquery
+    # criteria would otherwise make SQLAlchemy read back the deleted keys.
+    pakl_result = cast(
+        "CursorResult",
+        session.execute(
+            delete(PartitionedAssetKeyLog)
+            .where(
+                *pakl_where_clause,
+                exists().where(pakl_of_apdr, 
AssetPartitionDagRun.created_dag_run_id.is_(None)),
+            )
+            .execution_options(synchronize_session=False)
+        ),
+    )
+
+    apdr_where_clause: list[ColumnElement[bool]] = [
+        AssetPartitionDagRun.created_dag_run_id.is_(None),
+        ~exists().where(pakl_of_apdr),
+    ]
+    if dag_id is not None:
+        apdr_where_clause.append(AssetPartitionDagRun.target_dag_id == dag_id)
+    if permitted_dag_ids is not None:
+        
apdr_where_clause.append(AssetPartitionDagRun.target_dag_id.in_(permitted_dag_ids))
+    session.execute(
+        
delete(AssetPartitionDagRun).where(*apdr_where_clause).execution_options(synchronize_session=False)

Review Comment:
   One concurrency concern: two calls clearing different assets of the same 
partition could potentially leave the APDR empty and pending. Under READ 
COMMITTED, each call's `NOT EXISTS` can still see the other call's log rows 
until it commits, so neither call removes the APDR. I was able to reproduce 
this on Postgres 16 with two sessions.
   
   Would it make sense to wrap the candidate select from the note above with 
`with_row_locks(..., session=session, key_share=False)`? With that in place, 
the second call waits for the first, and the APDR is removed as expected.
   
   Once the cleanup is scoped to the candidate APDRs, I think these two changes 
would work well together, since a later call would no longer sweep an APDR left 
behind by the concurrent case.
   



##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py:
##########
@@ -124,6 +127,89 @@ def _generate_queued_event_where_clause(
     return where_clause
 
 
+QueuedEventPartitionKeyQuery = Annotated[
+    str | None,
+    Query(
+        description=(
+            "Delete queued events of partitioned assets emitted for this 
partition key whose "
+            "Dag run has not been created yet, instead of queued events of 
non-partitioned assets."
+        ),
+    ),
+]
+
+
+def _delete_pending_partitioned_queued_events(
+    *,
+    partition_key: str,
+    session: Session,
+    asset_id: int | None = None,
+    dag_id: str | None = None,
+    before: datetime | str | None = None,
+    permitted_dag_ids: set[str] | None = None,
+) -> int:
+    """
+    Delete queued partitioned asset events whose Dag run has not been created 
yet.
+
+    ``partition_key`` is matched against the partition key of the asset event, 
as in
+    ``POST /assets/{asset_id}/materialize``, not against the partition key of 
the Dag run
+    the event would create; a partition mapper can make the two differ.
+
+    Only events of an existing ``AssetPartitionDagRun`` whose Dag run has not 
been created
+    yet are deleted; ``PartitionedAssetKeyLog`` rows left behind once their 
Dag run was
+    deleted are history, not queued events. Pending ``AssetPartitionDagRun`` 
rows left
+    without any contributing events are deleted afterwards, so clearing some 
of the events
+    of a partition keeps the partition waiting for them rather than dropping 
the progress
+    made by the others.
+
+    :return: The number of deleted queued events.
+    """
+    pakl_where_clause: list[ColumnElement[bool]] = [
+        PartitionedAssetKeyLog.source_partition_key == partition_key,
+    ]
+    if asset_id is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.asset_id == asset_id)
+    if dag_id is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.target_dag_id == 
dag_id)
+    if before is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.created_at < before)
+    if permitted_dag_ids is not None:
+        
pakl_where_clause.append(PartitionedAssetKeyLog.target_dag_id.in_(permitted_dag_ids))
+    pakl_of_apdr = PartitionedAssetKeyLog.asset_partition_dag_run_id == 
AssetPartitionDagRun.id
+
+    # Nothing is loaded into the session beforehand, so skip synchronizing it; 
the subquery
+    # criteria would otherwise make SQLAlchemy read back the deleted keys.
+    pakl_result = cast(
+        "CursorResult",
+        session.execute(
+            delete(PartitionedAssetKeyLog)
+            .where(
+                *pakl_where_clause,
+                exists().where(pakl_of_apdr, 
AssetPartitionDagRun.created_dag_run_id.is_(None)),
+            )
+            .execution_options(synchronize_session=False)
+        ),
+    )
+
+    apdr_where_clause: list[ColumnElement[bool]] = [
+        AssetPartitionDagRun.created_dag_run_id.is_(None),
+        ~exists().where(pakl_of_apdr),
+    ]
+    if dag_id is not None:
+        apdr_where_clause.append(AssetPartitionDagRun.target_dag_id == dag_id)
+    if permitted_dag_ids is not None:
+        
apdr_where_clause.append(AssetPartitionDagRun.target_dag_id.in_(permitted_dag_ids))

Review Comment:
   Could we scope this cleanup to the APDRs touched by the first DELETE? It 
looks like the cleanup is currently only narrowed by `dag_id` / 
`permitted_dag_ids`, so `partition_key`, `asset_id`, and `before` don't affect 
which APDRs are removed.
   
   I noticed that clearing one asset's key through 
`/assets/{asset_id}/queuedEvents` can also remove an empty pending APDR 
belonging to another DAG, whereas 
`test_dag_endpoints_only_delete_events_queued_for_that_dag` keeps that APDR 
when using the DAG endpoints.
   
   Would it make sense to collect the candidate IDs first and then scope both 
DELETEs with `.in_(apdr_ids)`? Something along these lines:
   
   ```python
   apdr_ids = session.scalars(
       select(AssetPartitionDagRun.id).where(
           AssetPartitionDagRun.created_dag_run_id.is_(None),
           exists().where(pakl_of_apdr, *pakl_where_clause),
       )
   ).all()
   if not apdr_ids:
       return 0
   ```
   
   I tried this approach, and the 17 tests in 
`TestDeletePartitionedQueuedEvents` still pass for me.
   



##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py:
##########
@@ -124,6 +127,89 @@ def _generate_queued_event_where_clause(
     return where_clause
 
 
+QueuedEventPartitionKeyQuery = Annotated[
+    str | None,
+    Query(
+        description=(
+            "Delete queued events of partitioned assets emitted for this 
partition key whose "
+            "Dag run has not been created yet, instead of queued events of 
non-partitioned assets."
+        ),
+    ),
+]
+
+
+def _delete_pending_partitioned_queued_events(
+    *,
+    partition_key: str,
+    session: Session,
+    asset_id: int | None = None,
+    dag_id: str | None = None,
+    before: datetime | str | None = None,
+    permitted_dag_ids: set[str] | None = None,
+) -> int:
+    """
+    Delete queued partitioned asset events whose Dag run has not been created 
yet.
+
+    ``partition_key`` is matched against the partition key of the asset event, 
as in
+    ``POST /assets/{asset_id}/materialize``, not against the partition key of 
the Dag run
+    the event would create; a partition mapper can make the two differ.
+
+    Only events of an existing ``AssetPartitionDagRun`` whose Dag run has not 
been created
+    yet are deleted; ``PartitionedAssetKeyLog`` rows left behind once their 
Dag run was
+    deleted are history, not queued events. Pending ``AssetPartitionDagRun`` 
rows left
+    without any contributing events are deleted afterwards, so clearing some 
of the events
+    of a partition keeps the partition waiting for them rather than dropping 
the progress
+    made by the others.
+
+    :return: The number of deleted queued events.
+    """
+    pakl_where_clause: list[ColumnElement[bool]] = [
+        PartitionedAssetKeyLog.source_partition_key == partition_key,
+    ]
+    if asset_id is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.asset_id == asset_id)
+    if dag_id is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.target_dag_id == 
dag_id)
+    if before is not None:
+        pakl_where_clause.append(PartitionedAssetKeyLog.created_at < before)
+    if permitted_dag_ids is not None:
+        
pakl_where_clause.append(PartitionedAssetKeyLog.target_dag_id.in_(permitted_dag_ids))
+    pakl_of_apdr = PartitionedAssetKeyLog.asset_partition_dag_run_id == 
AssetPartitionDagRun.id
+
+    # Nothing is loaded into the session beforehand, so skip synchronizing it; 
the subquery
+    # criteria would otherwise make SQLAlchemy read back the deleted keys.
+    pakl_result = cast(
+        "CursorResult",
+        session.execute(
+            delete(PartitionedAssetKeyLog)
+            .where(
+                *pakl_where_clause,
+                exists().where(pakl_of_apdr, 
AssetPartitionDagRun.created_dag_run_id.is_(None)),
+            )
+            .execution_options(synchronize_session=False)
+        ),
+    )
+
+    apdr_where_clause: list[ColumnElement[bool]] = [
+        AssetPartitionDagRun.created_dag_run_id.is_(None),
+        ~exists().where(pakl_of_apdr),
+    ]
+    if dag_id is not None:
+        apdr_where_clause.append(AssetPartitionDagRun.target_dag_id == dag_id)
+    if permitted_dag_ids is not None:
+        
apdr_where_clause.append(AssetPartitionDagRun.target_dag_id.in_(permitted_dag_ids))
+    session.execute(
+        
delete(AssetPartitionDagRun).where(*apdr_where_clause).execution_options(synchronize_session=False)
+    )
+    return pakl_result.rowcount
+
+
+def _queued_event_not_found_detail(subject: str, partition_key: str | None) -> 
str:

Review Comment:
   Small naming nit: AGENTS.md recommends using action verbs for function 
names, and `_queued_event_not_found_detail` reads a bit like an attribute. 
Would `_build_queued_event_not_found_detail` be a better fit here? It has three 
call sites.
   



-- 
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