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

  1. Retry the pipeline launch; these failures are often transient network issues.
  2. Verify connectivity to the job server endpoint (--endpoint) and that it is running and healthy.
  3. Check for proxy/gRPC payload or timeout limits and increase deadlines if needed.
  4. 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

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


AI-assisted analysis of apache/beam@12126d8942 (2026-09-13). Data as JSON: /api/errors/00665b0cf059c136. Report an issue: GitHub.