GitHub user MentalOfCrow added a comment to the discussion: 
DataprocCreateBatchOperator: Retry on 5xx errors for deferrable

Hi Vit,

Yes. With a new `batch_id` on each Airflow task retry, the useful place to 
retry is the polling call inside **`DataprocBatchTrigger`**. That lets the 
current deferred attempt keep monitoring the batch it already submitted. It 
does not require changing your ID-generation policy.

There is one important distinction in the 15.1.0 code: 
[`DataprocAsyncHook.get_batch()`](https://github.com/apache/airflow/blob/providers-google/15.1.0/providers/google/src/airflow/providers/google/cloud/hooks/dataproc.py)
 already accepts `retry` and forwards it to the Google async client, with 
`DEFAULT` as its default. So the SDK may already retry some errors. However, 
[the 
trigger](https://github.com/apache/airflow/blob/providers-google/15.1.0/providers/google/src/airflow/providers/google/cloud/triggers/dataproc.py)
 has no recovery around an exception that escapes that call, and [the 
operator's deferral 
branch](https://github.com/apache/airflow/blob/providers-google/15.1.0/providers/google/src/airflow/providers/google/cloud/operators/dataproc.py)
 does not pass its `retry` setting to the trigger.

The metadata-plugin message also matters: it reports an `UNAVAILABLE` error 
while obtaining RPC metadata, rather than a returned Dataproc batch state of 
`FAILED`. Retrying a read of the same batch is appropriate for that transient 
error. The message alone does not identify which credential service produced 
the 503.

### A targeted implementation for 15.1.0

The same approach as the Dataflow/Cloud Run fixes you linked can be applied to 
Dataproc polling, but it needs a change in this trigger; those fixes are in 
other triggers.

For a **locally maintained provider patch**, add these module-level imports to 
`airflow/providers/google/cloud/triggers/dataproc.py` (`asyncio`, `Batch` and 
`TriggerEvent` are already imported there):

```python
import grpc
from grpc.aio import AioRpcError
from google.api_core.exceptions import ServiceUnavailable
```

Then replace `DataprocBatchTrigger.run()` with this illustrative version. It 
handles both the Google API-core `ServiceUnavailable` exception and the raw 
`AioRpcError` shown in your question, while checking the raw gRPC status 
explicitly:

```python
    async def run(self):
        consecutive_errors = 0
        max_consecutive_retries = 5  # Example policy; choose for your 
environment.
    
        while True:
            try:
                batch = await self.get_async_hook().get_batch(
                    project_id=self.project_id,
                    region=self.region,
                    batch_id=self.batch_id,
                )
            except (ServiceUnavailable, AioRpcError) as exc:
                if isinstance(exc, AioRpcError) and exc.code() != 
grpc.StatusCode.UNAVAILABLE:
                    raise
                consecutive_errors += 1
                if consecutive_errors > max_consecutive_retries:
                    raise
                delay = min(60.0, self.polling_interval_seconds * 2 ** 
(consecutive_errors - 1))
                self.log.warning(
                    "Transient UNAVAILABLE polling Dataproc; retry %s/%s in %s 
seconds",
                    consecutive_errors,
                    max_consecutive_retries,
                    delay,
                )
                await asyncio.sleep(delay)
                continue
    
            consecutive_errors = 0
            state = batch.state
            if state in (Batch.State.FAILED, Batch.State.SUCCEEDED, 
Batch.State.CANCELLED):
                break
            self.log.info("Current state is %s", state)
            await asyncio.sleep(self.polling_interval_seconds)
    
        yield TriggerEvent(
            {"batch_id": self.batch_id, "batch_state": state, 
"batch_state_message": batch.state_message}
        )
```

This method is adapted from the Apache-2.0-licensed provider 15.1.0 source 
linked above. Five retries and a 60-second maximum sleep are example policy 
values, not an upstream default. There are at most six consecutive failed 
polling calls before the error is raised. A successful status read resets that 
counter. The delay grows from your polling interval and is capped at 60 
seconds; for many concurrent triggers, add jitter to avoid synchronized retries.

The important properties are:

- Every retry uses the **same existing `batch_id`**. This code never calls 
`create_batch()`.
- Raw gRPC errors other than `UNAVAILABLE` are raised immediately. Ordinary 
permission, not-found and programming errors are not swallowed.
- A successfully fetched `FAILED` or `CANCELLED` batch still produces the 
original terminal event. It is not converted into success.
- `asyncio.sleep()` yields to the triggerer's event loop. Task cancellation is 
allowed to propagate.
- The **15.1.0 event payload stays unchanged**, including its `batch_state` 
value. Keep that contract when backporting; do not copy a newer trigger's event 
representation without checking the corresponding operator.

This bounds retries of errors that escape the SDK. It is not a hard wall-clock 
timeout for the whole batch or for each underlying RPC. If the outage outlasts 
the retry budget, the trigger still fails; any configured Airflow task retry 
can then create a new batch under your current ID policy.

### Applying and checking it

Apply this in a versioned provider build that you can maintain and roll back, 
then deploy it consistently to the relevant Airflow components, particularly 
the **triggerer**, which runs the polling code. Changing the DAG's `retry=` 
alone will not install this behavior. If you implement a separate custom 
trigger instead, its `serialize()` must return its own importable class path, 
and the operator must actually defer to that class; inheriting the built-in 
serialization unchanged would reload the built-in trigger.

For validation, inject one `UNAVAILABLE` polling error followed by `RUNNING`/ 
`SUCCEEDED`. Check that the task remains deferred during recovery, polling 
keeps the same batch name, and no second batch is submitted. Also check 
repeated outages past the budget, permission errors, and a genuinely failed 
batch.

I ran **eight isolated async tests** of the example method with a mocked hook 
and event/state objects, including real `grpc.aio.AioRpcError` objects. They 
covered recovery, the retry limit, counter reset, non-retryable errors, 
failed/cancelled batch events and cancellation. The API-core exception was a 
stand-in in that harness. I have **not** run this inside Airflow or against 
Google Cloud, so it still needs integration testing with your installed 
dependency versions and deployment.

### Upgrade or temporary alternative

I also checked [the current `main` trigger 
source](https://github.com/apache/airflow/blob/main/providers/google/src/airflow/providers/google/cloud/triggers/dataproc.py):
 it still has an unguarded polling call. I would therefore not promise that an 
arbitrary provider upgrade fixes this particular failure. A newer SDK could 
affect the underlying RPC behavior, but that is a separate question from adding 
this trigger-level recovery.

If your environment does not permit a provider/custom-trigger change and you 
can afford to occupy a worker while waiting, `deferrable=False` uses the 
synchronous `wait_for_batch()` path. In 15.1.0, that path already catches 
`ServerError` while polling and forwards the operator's `retry` setting. That 
is a temporary alternative to investigate, with different resource and failure 
behavior—not proof that the same credential failure will disappear.

For an upstream fix, the focused scope would be the Dataproc trigger's polling 
recovery plus regression tests, with the retry policy and event compatibility 
reviewed by maintainers. Your two linked PRs are useful precedents; they do not 
by themselves change Dataproc's implementation.

*AI-assisted answer; source inspection and isolated-test limits are described 
above.*

GitHub link: 
https://github.com/apache/airflow/discussions/73675#discussioncomment-18712255

----
This is an automatically sent email for [email protected].
To unsubscribe, please send an email to: [email protected]

Reply via email to