apache / apache/beam

Improve the recoverability of DoFn

Open
#27,094 3 comments 0 reactions 0 assignees View on GitHub
awaiting triage java new feature P2
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

### What would you like to happen?

During a recent Dataflow streaming job, an ongoing DoFn experienced a period of frequent 429 error responses from an overloaded Elasticsearch sink (using ElasticsearchIO). During this time, it was observed that autoscaling occurred due to increased backlog, but workers did not re-establish their connections once the sink was back online, causing the pipeline to remain frozen.

Steps should be taken to improve the resiliency of DoFn to these sorts of disruptions and ensure that they resolve themselves without the pipeline needing to be restarted manually.

### Issue Priority

Priority: 2 (default / most feature requests should be filed as P2)

### Issue Components

- [ ] Component: Python SDK
- [X] Component: Java SDK
- [ ] Component: Go SDK
- [ ] Component: Typescript SDK
- [ ] Component: IO connector
- [ ] Component: Beam examples
- [ ] Component: Beam playground
- [ ] Component: Beam katas
- [ ] Component: Website
- [ ] Component: Spark Runner
- [ ] Component: Flink Runner
- [ ] Component: Samza Runner
- [ ] Component: Twister2 Runner
- [ ] Component: Hazelcast Jet Runner
- [ ] Component: Google Cloud Dataflow Runner

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.