Apache Beam Go SDK grades示例在Spark 2.4.5集群运行时出现容器异常及文件缺失报错求助
I'm trying to run the Apache Beam Go SDK's grades example using the Spark Runner on a 1-master 2-slave Spark 2.4.5 cluster. Both SSH and Docker are installed and running, but I keep hitting errors that I can't pin down the root cause of.
Commands Used
I started the job service endpoint with this command:
docker run --net=host apache/beam_spark_job_server:latest --spark-master-url=spark://master:7077 --artifacts-dir /tmp/beam-artifact-staging
And ran the example with:
grades -runner=spark -endpoint=localhost:8099 -job_name=gradetest
Error Messages
Driver Error Output
Failed to retrieve staged files: failed to retrieve /tmp/staged in 3 attempts: failed to retrieve chunk for /tmp/staged/worker caused by: rpc error: code = Unknown desc = ; failed to retrieve chunk for /tmp/staged/worker caused by: rpc error: code = Unknown desc = ; failed to retrieve chunk for /tmp/staged/worker caused by: rpc error: code = Unknown desc = ; failed to retrieve chunk for /tmp/staged/worker caused by: rpc error: code = Unknown desc = 21/09/19 11:01:47 WARN BlockManager: Putting block rdd_2_1 failed due to exception org.apache.beam.vendor.guava.v26_0_jre.com.google.common.util.concurrent.UncheckedExecutionException: java.lang.IllegalStateException: No container running for id Driver commanded a shutdown apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2050) at org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.get(LocalCache.java:3952) at org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.getOrLoad(LocalCache.java:3974) at
Spark Standard Error Log
23/09/20 10:06:22 INFO Executor: Adding file:/opt/spark/work/app-20210920100619-0017/1/./beam-runners-spark-job-server.jar to class loader 23/09/20 10:06:22 INFO TorrentBroadcast: Started reading broadcast variable 0 23/09/20 10:06:22 INFO TransportClientFactory: Successfully created connection to master/192.168.1.*:44365 after 1 ms (0 ms spent in bootstraps) 23/09/20 10:06:22 INFO MemoryStore: Block broadcast_0_piece0 stored as bytes in memory (estimated size 10.9 KB, free 366.3 MB) 23/09/20 10:06:22 INFO TorrentBroadcast: Reading broadcast variable 0 took 97 ms 23/09/20 10:06:22 INFO MemoryStore: Block broadcast_0 stored as values in memory (estimated size 24.9 KB, free 366.3 MB) java.io.FileNotFoundException: /tmp/beam-artifact-staging/e40099113cf8136935edc839aa85487c0532034c0a63f8cbadd7fccac0f98ed0/1-go-worker (No such file or directory) at java.io.FileInputStream.open0(Native Method) at java.io.FileInputStream.open(FileInputStream.java:195) at java.io.FileInputStream.<init>(FileInputStream.java:138) at org.apache.beam.sdk.io.LocalFileSystem.open(LocalFileSystem.java:127) at org.apache.beam.sdk.io.LocalFileSystem.open(LocalFileSystem.java:83) at org.apache.beam.sdk.io.FileSystems.open(FileSystems.java:257) at org.apache.beam.runners.fnexecution.artifact.ArtifactRetrievalService.getArtifact(ArtifactRetrievalService.java:124) at
Root Cause Analysis
The critical error here is the FileNotFoundException for the 1-go-worker binary in the Spark executor logs. This tells us the Go worker binary (required for Beam Go SDK jobs) isn't being properly shared across all nodes in your Spark cluster.
When you run the Beam Spark Job Server in Docker with --net=host, the /tmp/beam-artifact-staging directory is only accessible on the master node where the container is running. Spark slave nodes can't reach this local directory, so when they try to fetch the worker binary, it's missing. The driver's "failed to retrieve staged files" error confirms this artifact distribution breakdown.
Step-by-Step Solutions
1. Use a Shared Network File System (NFS) for Artifact Staging
All cluster nodes need access to the same artifact staging directory. Set up an NFS share and mount it consistently on master and slaves (e.g., /mnt/beam-artifacts):
- Create the NFS share on the master and configure exports to allow slave access
- Mount the share on each slave node (use
mounttemporarily or add to/etc/fstabfor persistence) - Update your job server command to use this shared path:
docker run --net=host -v /mnt/beam-artifacts:/tmp/beam-artifact-staging apache/beam_spark_job_server:latest --spark-master-url=spark://master:7077 --artifacts-dir /tmp/beam-artifact-staging
2. Verify Go Worker Binary Staging
After starting the job, check the shared staging directory on the master to confirm the 1-go-worker binary exists. If it's missing, ensure the Beam Go SDK is correctly building the worker during job submission.
3. Fix Spark Worker Permissions
Make sure the user running Spark workers has read access to the shared artifact directory:
- On a slave node, run
ls -l /mnt/beam-artifacts/<artifact-id>/1-go-workerto confirm accessibility - Adjust permissions with
chmodor add the Spark worker user to the directory's group if needed
4. Match Beam and Spark Versions
The latest Beam job server image might have compatibility issues with Spark 2.4.5. Try a Beam version known to work with Spark 2.4.5 (e.g., Beam 2.30.0):
docker run --net=host -v /mnt/beam-artifacts:/tmp/beam-artifact-staging apache/beam_spark_job_server:2.30.0 --spark-master-url=spark://master:7077 --artifacts-dir /tmp/beam-artifact-staging
5. Use Spark's Distributed Cache (Testing Only)
For quick testing, bypass the shared directory by using Spark's built-in distributed cache to stage artifacts. Add this flag when running the example:
grades -runner=spark -endpoint=localhost:8099 -job_name=gradetest -spark.use_remote_artifacts=true
内容的提问来源于stack exchange,提问作者lpad ze

