AppflowHook.run_flow raises KeyError when the execution record is not on the first page of DescribeFlowExecutionRecords
- Dominant language
- Python
- Stars
- 46.9k
- Forks
- 17.8k
- Avg merge
- 2d 9h
- Merged PRs (30d)
- 472
Description
### Under which category would you file this issue?
Providers
### Apache Airflow version
3.3.1
### What happened and how to reproduce it?
`AppflowHook.run_flow` finishes by calling `_log_execution_description`, which looks up the execution it just ran in the output of `DescribeFlowExecutionRecords`:
```python
def _log_execution_description(self, flow_name: str, execution_id: str):
response_desc = self.conn.describe_flow_execution_records(flowName=flow_name)
last_execs = {fe["executionId"]: fe for fe in response_desc["flowExecutions"]}
exec_details = last_execs[execution_id]
self.log.info("Run complete, execution details: %s", exec_details)
```
That is one call with no `maxResults` and no `nextToken` handling, so it only ever sees the first page. The AppFlow API model says `maxResults` defaults to 20 and the response carries a `nextToken`, and nothing in the docs says the records come back newest first. So the execution you just started is not guaranteed to be in that first page. When it is not, `last_execs[execution_id]` raises `KeyError` and the task fails even though AppFlow reported the flow as complete.
`run_flow` defaults to `wait_for_completion=True`, and every AppFlow operator goes through it, so this is on the normal path for `AppflowRunOperator` and friends.
Worth noting `AppflowRecordsShortCircuitOperator._has_new_records_func` already treats this same API as paginated, asking for 100 records and looping on `nextToken`. The hook just does not.
**Steps to reproduce**
No AWS account needed:
```python
from unittest import mock
from airflow.providers.amazon.aws.hooks.appflow import AppflowHook
FLOW = "flow0"
TARGET = "ex_target"
# AWS returns at most maxResults records (default 20) plus a nextToken.
page1 = {
"flowExecutions": [
{"executionId": f"other_{i}", "executionStatus": "Successful"} for i in range(20)
],
"nextToken": "page2",
}
page2 = {
"flowExecutions": [
{"executionId": TARGET, "executionStatus": "Successful",
"executionResult": {"recordsProcessed": 42}}
]
}
with mock.patch.object(AppflowHook, "conn") as conn:
conn.start_flow.return_value = {"executionId": TARGET}
conn.describe_flow_execution_records.side_effect = (
lambda **kw: page2 if kw.get("nextToken") == "page2" else page1
)
hook = AppflowHook(aws_conn_id="aws_default", region_name="us-east-1")
with mock.patch("airflow.providers.amazon.aws.waiters.base_waiter.BaseBotoWaiter.waiter"):
hook.run_flow(flow_name=FLOW, poll_interval=0)
```
Observed:
```
[info] executionId: ex_target
Traceback (most recent call last):
...
File ".../airflow/providers/amazon/aws/hooks/appflow.py", line 86, in _log_execution_description
exec_details = last_execs[execution_id]
~~~~~~~~~~^^^^^^^^^^^^^^
KeyError: 'ex_target'
```
`describe_flow_execution_records` is called exactly once, as `call(flowName='flow0')`.
### What you think should happen instead?
The lookup should page through `DescribeFlowExecutionRecords` until it finds the execution, and should ask for the maximum page size so that is normally one call.
If the record genuinely is not there, it should log a warning rather than raise. The whole job of `_log_execution_description` is to write one log line, so it should not be able to fail a run that AppFlow already reported as successful.
### Operating System
Linux (reproduced in the `python:3.12-slim` container, kernel 6.6.87.2-microsoft-standard-WSL2)
### Deployment
Other
### Deployment details
Reproduced against released packages in a clean container, not a full Airflow deployment. The failing path is pure hook logic with the boto3 client mocked.
### Apache Airflow Provider(s)
amazon
### Versions of Apache Airflow Providers
```
apache-airflow==3.3.1
apache-airflow-providers-amazon==9.35.1
botocore==1.43.90
Python 3.12.14
```
### Anything else?
The waiter in `run_flow` has the same blind spot. Its acceptors filter `flowExecutions[?executionId=='{{EXECUTION_ID}}']` out of the same unpaginated call, so with more than 20 records the waiter can poll until `max_attempts` runs out on a flow that actually finished. Botocore waiters cannot paginate, so that one cannot be fully fixed the same way, but the waiter can at least ask for 100 records per page instead of 20.
### Are you willing to submit PR?
- [X] Yes I am willing to submit a PR!
### Code of Conduct
- [X] I agree to follow this project's Code of Conduct
Contributor guide
Research direction
Start in airflow/providers/amazon/aws/hooks/appflow.py with _log_execution_description and the waiter used by run_flow; run the supplied mocked reproduction first. Trace the DescribeFlowExecutionRecords calls and nextToken handling. Done means the execution can be found across pages, a missing record only logs a warning, and the waiter requests the larger page size without failing a completed run.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- aws, python
- Domain
- backend, cloud
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Active
- Clarity
- Clearly specified
- Newbie friendliness
- 68/100