【问题标题】:java.lang.RuntimeException: com.datastax.bdp.fs.model.NoSuchFileException: File not found: /tmp/hive/java.lang.RuntimeException:com.datastax.bdp.fs.model.NoSuchFileException:找不到文件:/tmp/hive/
【发布时间】:2018-10-30 21:15:58
【问题描述】:

我有以下代码:

def main(args: Array[String]) {

val conf = new SparkConf()
  .setAppName("Fleet")
  .set("spark.executor.memory", "1g")
  .set("spark.driver.memory", "2g")
  .set("spark.submit.deployMode", "cluster")
  .set("spark.executor.instances", "4")
  .set("spark.executor.cores", "3")
  .set("spark.cores.max", "12")
  .set("spark.driver.cores", "4")
  .set("spark.ui.port", "4040")
  .set("spark.streaming.backpressure.enabled", "true")
  .set("spark.streaming.kafka.maxRatePerPartition", "30")

val spark = SparkSession
  .builder
  .appName("Fleet")
  .config("spark.cassandra.connection.host", "192.168.0.40")
  .config("spark.cassandra.connection.port", "9042")
  .config("spark.submit.deployMode", "cluster")
  .master("local[*]")
  .getOrCreate()

val sc = SparkContext.getOrCreate(conf)
val ssc = new StreamingContext(sc, Seconds(10))
val sqlContext = new SQLContext(sc)
val topics = Map("historyfleet" -> 1) 
val kafkaStream = KafkaUtils.createStream(ssc, "192.168.0.40:2181", "fleetgroup", topics)

kafkaStream.foreachRDD(rdd =>
  {
    val dfs = rdd.toDF()
    println(dfs.show())
    dfs.write.format("org.apache.spark.sql.cassandra").options(Map("table" -> "test", "keyspace" -> "test_db")).mode(SaveMode.Append).save()
  })
ssc.start()
ssc.awaitTermination()

}

我可以在本地机器上从 Eclipse 执行这个程序,但是当尝试通过集群上的 spark 提交作业执行时,它给出了一个错误:-

ERROR 2018-05-21 13:00:27,009 org.apache.spark.deploy.DseSparkSubmitBootstrapper: Failed to start or submit Spark application
java.lang.RuntimeException: com.datastax.bdp.fs.model.NoSuchFileException: File not found: /tmp/hive/
at org.apache.hadoop.hive.ql.session.SessionState.start(SessionState.java:522) ~[hive-exec-1.2.1.spark2.jar:1.2.1.spark2]
at org.apache.spark.sql.hive.client.HiveClientImpl.<init>(HiveClientImpl.scala:189) ~[spark-hive_2.11-2.0.2.16.jar:2.0.2.16]
at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method) ~[na:1.8.0_161]
at sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62) ~[na:1.8.0_161]
at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45) ~[na:1.8.0_161]
at java.lang.reflect.Constructor.newInstance(Constructor.java:423) ~[na:1.8.0_161]
at org.apache.spark.sql.hive.client.IsolatedClientLoader.createClient(IsolatedClientLoader.scala:258) ~[spark-hive_2.11-2.0.2.16.jar:2.0.2.16]
at org.apache.spark.sql.hive.HiveUtils$.newClientForMetadata(HiveUtils.scala:359) ~[spark-hive_2.11-2.0.2.16.jar:2.0.2.16]

我的想法是从 Kafka 流中获取记录并将数据推送到 Cassandra。谢谢,

【问题讨论】:

  • 您是否在 Windows 集群中运行?也摆脱master("local[*]"),但我相信与错误无关
  • 不,它是一个 Linux 集群。并且每当我在集群上部署我的应用程序时,总是评论 - master("local[*]") this。以确保它应该部署在集群上。

标签: scala apache-spark dataframe datastax-enterprise spark-cassandra-connector


【解决方案1】:

您需要将 dsefs 键空间的复制因子 (rf) 增加到大于 1 的值。此外,dsefs 键空间(与任何其他键空间一样)最适用于 NetworkTopologyStrategy。这是一个使用 rf = 3 更改策略的命令。

ALTER KEYSPACE dsefs WITH replication = {'class':'NetworkTopologyStrategy', '<YOUR DC HERE>': '3'}

更改键空间后,您需要运行节点修复所有个节点。

nodetool repair dsefs

除此之外,您可以从 DSEFS 中删除 /tmp/hive 并使用重新创建它

dse fs
mkdir -p -m 733 /tmp/hive

【讨论】:

  • 你确定吗? dsefs 是文件系统键空间。我可能会失去整个集群
  • NetworkTopologyStrategy 是 DSEFS 文档推荐的 docs.datastax.com/en/dse/5.1/dse-dev/datastax_enterprise/…
  • 感谢您的帮助。我将使用 dummy cluster 尝试此解决方案。现在我的问题得到了解决。我在正常模式下启动 Cassandra。我必须从 -k 开始启用分析模式。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-11-15
  • 2023-04-10
  • 2018-11-04
  • 2012-03-05
  • 2017-03-11
  • 1970-01-01
相关资源
最近更新 更多