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]
