apache/beam · error
StageFile chunk send failed
Error message
StageFile chunk send failed
What it means
While streaming 1MB chunks of a staged artifact to the job server, stream.Send returned an error other than io.EOF, wrapped as 'StageFile chunk send failed'. This means the gRPC reverse artifact retrieval stream broke mid-upload.
Source
Thrown at sdks/go/pkg/beam/runners/universal/runnerlib/stage.go:165
}
defer fd.Close()
data := make([]byte, 1<<20)
for {
n, err := fd.Read(data)
if n > 0 {
sendErr := stream.Send(&jobpb.ArtifactResponseWrapper{
Response: &jobpb.ArtifactResponseWrapper_GetArtifactResponse{
GetArtifactResponse: &jobpb.GetArtifactResponse{
Data: data[:n],
},
}})
if sendErr == io.EOF {
return sendErr
}
if sendErr != nil {
return errors.Wrap(sendErr, "StageFile chunk send failed")
}
}
if err == io.EOF {
sendErr := stream.Send(&jobpb.ArtifactResponseWrapper{
IsLast: true,
Response: &jobpb.ArtifactResponseWrapper_GetArtifactResponse{
GetArtifactResponse: &jobpb.GetArtifactResponse{},
}})
return sendErr
}
if err != nil {
return err
}
}
}
View on GitHub (pinned to 12126d8942)
Solutions
- Retry the pipeline launch; these failures are often transient network issues.
- Verify connectivity to the job server endpoint (--endpoint) and that it is running and healthy.
- Check for proxy/gRPC payload or timeout limits and increase deadlines if needed.
- Inspect job server logs for stream cancellation or resource limits.
Example fix
// before
if sendErr != nil {
return errors.Wrap(sendErr, "StageFile chunk send failed")
}
// after
if sendErr != nil {
if status.Code(sendErr) == codes.Unavailable {
return retryable(sendErr) // caller retries with backoff
}
return errors.Wrap(sendErr, "StageFile chunk send failed")
} Defensive patterns
Strategy: retry
Validate before calling
conn, err := grpc.Dial(endpoint, grpc.WithBlock(), grpc.WithTimeout(5*time.Second))
if err != nil {
return fmt.Errorf("job server %s unreachable: %w", endpoint, err)
} Try / catch
err := stageFiles(...)
if err != nil {
if status.Code(errors.Cause(err)) == codes.Unavailable || status.Code(errors.Cause(err)) == codes.DeadlineExceeded {
// retry staging with exponential backoff
}
} Prevention
- Check job server health before launching.
- Avoid staging very large artifacts over unstable networks.
- Set generous gRPC deadlines for artifact staging.
When it happens
Trigger: The gRPC connection to the job server drops, times out, or the server cancels the RPC while sending artifact chunks; any non-EOF send error during the upload loop.
Common situations: Uploading large artifacts over an unstable network; job server restarted or deadline exceeded during staging; proxy or firewall killing long-lived gRPC streams.
Related errors
- failed to send chunks for %v; close error: %v
- chunk send failed
- failed to connect to state service %v
- StateChannel[%v].Send(%v): channel closed due to: %v
- failed to connect: %v
AI-assisted analysis of apache/beam@12126d8942 (2026-09-13).
Data as JSON: /api/errors/00665b0cf059c136.
Report an issue: GitHub.