Ошибка при запуске spark-submit как yarn и cluster
Проект нормально открывается только при deploy-mode client, но при переключении на cluster, во время деплоя проекта выходит ошибка
2020-10-16 04:37:56,600 INFO yarn.Client: Application report for application_1602820124289_0008 (state: ACCEPTED)
2020-10-16 04:37:57,602 INFO yarn.Client: Application report for application_1602820124289_0008 (state: ACCEPTED)
2020-10-16 04:37:58,605 INFO yarn.Client: Application report for application_1602820124289_0008 (state: FAILED)
2020-10-16 04:37:58,605 INFO yarn.Client:
client token: N/A
diagnostics: Application application_1602820124289_0008 failed 2 times due to AM Container for appattempt_1602820124289_0008_000002 exited with exitCode: 13
Failing this attempt.Diagnostics: [2020-10-16 04:37:57.600]Exception from container-launch.
Container id: container_1602820124289_0008_02_000001
Exit code: 13
[2020-10-16 04:37:57.602]Container exited with a non-zero exit code 13. Error file: prelaunch.err.
Last 4096 bytes of prelaunch.err :
Last 4096 bytes of stderr :
ionMaster.scala:516)
at org.apache.spark.deploy.yarn.ApplicationMaster.run(ApplicationMaster.scala:264)
at org.apache.spark.deploy.yarn.ApplicationMaster$$anon$3.run(ApplicationMaster.scala:890)
at org.apache.spark.deploy.yarn.ApplicationMaster$$anon$3.run(ApplicationMaster.scala:889)
at java.security.AccessController.doPrivileged(Native Method)
at javax.security.auth.Subject.doAs(Subject.java:422)
at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1746)
at org.apache.spark.deploy.yarn.ApplicationMaster$.main(ApplicationMaster.scala:889)
at org.apache.spark.deploy.yarn.ApplicationMaster.main(ApplicationMaster.scala)
2020-10-16 04:37:57,246 INFO spark.SparkContext: Invoking stop() from shutdown hook
2020-10-16 04:37:57,254 INFO server.AbstractConnector: Stopped Spark@340f8b2c{HTTP/1.1,[http/1.1]}{0.0.0.0:0}
2020-10-16 04:37:57,256 INFO ui.SparkUI: Stopped Spark web UI at http://10.110.160.124:43761
2020-10-16 04:37:57,280 INFO spark.MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped!
2020-10-16 04:37:57,294 INFO memory.MemoryStore: MemoryStore cleared
2020-10-16 04:37:57,294 INFO storage.BlockManager: BlockManager stopped
2020-10-16 04:37:57,305 INFO storage.BlockManagerMaster: BlockManagerMaster stopped
2020-10-16 04:37:57,310 INFO scheduler.OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped!
2020-10-16 04:37:57,321 INFO spark.SparkContext: Successfully stopped SparkContext
2020-10-16 04:37:57,322 INFO yarn.ApplicationMaster: Deleting staging directory hdfs://10.110.160.124:9000/user/hdoop/.sparkStaging/application_1602820124289_0008
2020-10-16 04:37:57,357 INFO concurrent.ThreadPoolTaskExecutor: Shutting down ExecutorService 'applicationTaskExecutor'
2020-10-16 04:37:57,357 INFO spark.SparkContext: SparkContext already stopped.
2020-10-16 04:37:57,546 ERROR util.Utils: Uncaught exception in thread Thread-2
java.lang.IllegalStateException: Shutdown in progress, cannot add a shutdownHook
at org.apache.hadoop.util.ShutdownHookManager.addShutdownHook(ShutdownHookManager.java:152)
at org.apache.hadoop.tracing.SpanReceiverHost.get(SpanReceiverHost.java:79)
at org.apache.hadoop.hdfs.DFSClient.<init>(DFSClient.java:634)
at org.apache.hadoop.hdfs.DFSClient.<init>(DFSClient.java:619)
at org.apache.hadoop.hdfs.DistributedFileSystem.initialize(DistributedFileSystem.java:149)
at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:2669)
at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:94)
at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:2703)
at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:2685)
at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:373)
at org.apache.hadoop.fs.Path.getFileSystem(Path.java:295)
at org.apache.spark.deploy.yarn.ApplicationMaster.cleanupStagingDir(ApplicationMaster.scala:674)
at org.apache.spark.deploy.yarn.ApplicationMaster.$anonfun$run$2(ApplicationMaster.scala:258)
at org.apache.spark.util.SparkShutdownHook.run(ShutdownHookManager.scala:214)
at org.apache.spark.util.SparkShutdownHookManager.$anonfun$runAll$2(ShutdownHookManager.scala:188)
at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1932)
at org.apache.spark.util.SparkShutdownHookManager.$anonfun$runAll$1(ShutdownHookManager.scala:188)
at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
at scala.util.Try$.apply(Try.scala:213)
at org.apache.spark.util.SparkShutdownHookManager.runAll(ShutdownHookManager.scala:188)
at org.apache.spark.util.SparkShutdownHookManager$$anon$2.run(ShutdownHookManager.scala:178)
at org.apache.hadoop.util.ShutdownHookManager$1.run(ShutdownHookManager.java:54)
2020-10-16 04:37:57,546 INFO util.ShutdownHookManager: Shutdown hook called
2020-10-16 04:37:57,547 INFO util.ShutdownHookManager: Deleting directory /home/hdoop/tmpdata/nm-local-dir/usercache/hdoop/appcache/application_1602820124289_0008/spark-c6a978d5-7f70-477f-9a2c-dc57cd583fe2
Вот мой сам проект
@Configuration
@EnableAutoConfiguration(exclude = {org.springframework.boot.autoconfigure.gson.GsonAutoConfiguration.class})
public class SparkConfig {
@Value("Test Spark Application")
private String appName;
@Value("local[*]")
private String masterUri;
@Bean
public SparkConf sparkConf() {
return new SparkConf()
.setAppName(appName)
.setMaster(masterUri)
.set("spark.sql.debug.maxToStringFields", "1000")
.setJars(new String[]
{"/home/hdoop/SparkApplication/demo-0.0.1-SNAPSHOT.jar"
, "/home/hdoop/SparkApplication/spark-core_2.12-3.0.1.jar"
, "/home/hdoop/SparkApplication/postgresql-42.2.10.jar"});
}
@Bean
public SparkContext sparkContext() {
return SparkContext.getOrCreate(sparkConf());
}
@Bean
public SparkSession sparkSession() {
return SparkSession
.builder()
.sparkContext(sparkContext())
.getOrCreate();
}
@Service
@EnableAutoConfiguration(exclude = {org.springframework.boot.autoconfigure.gson.GsonAutoConfiguration.class})
public class StackOverFlow implements Serializable {
@Autowired
private SparkConf sparkConf;
@Autowired
private SparkContext sparkContext;
@Autowired
private SparkSession session;
public List<Order> getObject(String param, String value, Long limit) {
Encoder<Order> orderEncoder = Encoders.bean(Order.class);
if (sparkContext.isStopped()){
sparkContext = SparkContext.getOrCreate(sparkConf);
}
SparkSession session = SparkSession
.builder()
.sparkContext(sparkContext)
.enableHiveSupport()
.getOrCreate();
SparkSession.setActiveSession(session);
DataFrameReader dfr = session.read();
Properties properties = new Properties();
properties.setProperty("driver", "org.postgresql.Driver");
properties.setProperty("user", "root");
properties.setProperty("password", "password");
properties.setProperty("query", "select * from orders o where " + param + " = '" + value + "' limit " + limit);
Dataset<Row> orderDataset = dfr
.jdbc(
"jdbc:postgresql://postgres:5432/orders",
"orders",
properties);
List<Order> orders = orderDataset.as(orderEncoder).collectAsList();
sparkContext.stop();
session.stop();
session.close();
return orders;
}