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]