DataprocCreateBatchOperator: Retry on 5xx errors for deferrable #73675
Replies: 3 comments
|
You're right that the trigger doesn't retry. The error itself ( The practical fix is to let Airflow's task retry handle it, and make the retry re-attach to the same batch instead of creating a new one. The operator already supports this: if DataprocCreateBatchOperator(
task_id="run_batch",
batch_id="my-job-{{ ds_nodash }}", # stable across retries, unique per run
batch={...},
region="europe-west1",
deferrable=True,
retries=3,
retry_delay=timedelta(minutes=1),
retry_exponential_backoff=True,
)With this, a 503 in the triggerer fails the attempt, Airflow retries with exponential backoff, the new attempt hits Two things to keep in mind:
If you want the trigger itself to retry transient errors, that would be a feature request for the Google provider. |
|
Hello, thank you for your answer and suggestion. We currently generate unique batch_id for each task run, even for retries, so this would not work for us. The same issue was already solved for Dataflow and Cloud Run. Would it be possible to implement the same retry logic in DataprocCreateBatchOperator? Best regards, |
|
Hi Vit, Yes. With a new There is one important distinction in the 15.1.0 code: The metadata-plugin message also matters: it reports an A targeted implementation for 15.1.0The 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 import grpc
from grpc.aio import AioRpcError
from google.api_core.exceptions import ServiceUnavailableThen replace 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:
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 itApply 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 For validation, inject one I ran eight isolated async tests of the example method with a mocked hook and event/state objects, including real Upgrade or temporary alternativeI also checked the current If your environment does not permit a provider/custom-trigger change and you can afford to occupy a worker while waiting, 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. |
Uh oh!
There was an error while loading. Please reload this page.
Hello,
we are using Airflow 2.11, apache-airflow-providers-google 15.1.0.
We run Dataproc Serverless batches via DataprocCreateBatchOperator with defferable = True. We are getting randomly error:
grpc.aio._call.AioRpcError: <AioRpcError of RPC that terminated with: status = StatusCode.UNAVAILABLE details = "Getting metadata from plugin failed with error: ("Error code {'code': 503, 'message': 'The service is currently unavailable.', 'status': 'UNAVAILABLE'}", '{\n "error": {\n "code": 503,\n "message": "The service is currently unavailable.",\n "status": "UNAVAILABLE"\n }\n}\n')" debug_error_string = "UNKNOWN:Error received from peer {grpc_message:"Getting metadata from plugin failed with error: (\"Error code {\'code\': 503, \'message\': \'The service is currently unavailable.\', \'status\': \'UNAVAILABLE\'}\", \'{\\n \"error\": {\\n \"code\": 503,\\n \"message\": \"The service is currently unavailable.\",\\n \"status\": \"UNAVAILABLE\"\\n }\\n}\\n\')", grpc_status:14, created_time:"2026-09-24T10:12:11.505328013+02:00"}"When I look at params of call DataprocBatchTrigger, there are no retry options passed there.
According to GCP docs:
Any client interacting with Google Cloud APIs must implement retry logic with exponential backoff for 5xx errors (specifically 500, 503 and 409")"Could somebody please help us how to solve this? Would update to latest version of providers help?
Thank you for help.
All reactions