小编eli*_*sah的帖子

Spark Standalone Mode多个shell会话(应用程序)

在具有多个工作节点的Spark 1.0.0独立模式中,我正在尝试从两台不同的计算机(同一Linux用户)运行Spark shell.

在文档中,它说"默认情况下,提交给独立模式群集的应用程序将以FIFO(先进先出)顺序运行,每个应用程序将尝试使用所有可用节点."

每个工作程序的核心数设置为4,其中8个可用(通过SPARK_JAVA_OPTS =" - Dspark.cores.max = 4").内存也是有限的,因此两者都应该可用.

但是,在查看Spark Master WebUI时,稍后启动的shell应用程序将始终保持"WAITING"状态,直到退出第一个.分配给它的核心数是0,每个节点10G的内存(与已经运行的核心相同)

有没有办法让两个shell同时运行而不使用Mesos?

apache-spark

6
推荐指数
1
解决办法
3330
查看次数

Spark源代码:如何理解withScope方法

我无法理解withScope方法的功能(实际上,我真的不知道RDDOperationScope类的含义)

特别是,withScope方法的参数列表中(body:=> T)的含义是什么:

private[spark] def withScope[T](
  sc: SparkContext,
  name: String,
  allowNesting: Boolean,
  ignoreParent: Boolean)(body: => T): T = {
// Save the old scope to restore it later
val scopeKey = SparkContext.RDD_SCOPE_KEY
val noOverrideKey = SparkContext.RDD_SCOPE_NO_OVERRIDE_KEY
val oldScopeJson = sc.getLocalProperty(scopeKey)
val oldScope = Option(oldScopeJson).map(RDDOperationScope.fromJson)
val oldNoOverride = sc.getLocalProperty(noOverrideKey)
try {
  if (ignoreParent) {
    // Ignore all parent settings and scopes and start afresh with our own root scope
    sc.setLocalProperty(scopeKey, new RDDOperationScope(name).toJson)
  } else if (sc.getLocalProperty(noOverrideKey) == null) {
    // Otherwise, …
Run Code Online (Sandbox Code Playgroud)

scala apache-spark

6
推荐指数
1
解决办法
1776
查看次数

Spark数据帧使用随机数据添加新列

我想在数据框中添加一个新列,其值由0或1组成.我使用了'randint'函数,

from random import randint

df1 = df.withColumn('isVal',randint(0,1))
Run Code Online (Sandbox Code Playgroud)

但我得到以下错误,

/spark/python/pyspark/sql/dataframe.py",第1313行,在withColumn断言isinstance(col,Column)中,"col应该是列"AssertionError:col应该是Column

如何使用自定义函数或randint函数为列生成随机值?

python apache-spark apache-spark-sql pyspark

6
推荐指数
2
解决办法
2万
查看次数

Apache Spark - 为什么要删除执行程序?“空闲”是什么意思?

使用动态资源分配

当执行程序空闲超过 spark.dynamicAllocation.executorIdleTimeout 秒时,Spark 应用程序会删除执行程序。

当我按如下方式设置执行程序空闲超时属性时spark.dynamicAllocation.executorIdleTimeout= 300,它会抛出以下警告

spark.ExecutorAllocationManager: Removing executor 0 because it has been idle for 300 seconds (new desired total will be 2)
Run Code Online (Sandbox Code Playgroud)

“空闲”是什么意思?这是否意味着工人不使用 CPU?对数据库的阻塞调用算作空闲吗?

apache-spark

6
推荐指数
1
解决办法
3581
查看次数

何时使用SPARK_CLASSPATH或SparkContext.addJar

我正在使用一个独立的火花群,一个主人和两个工人.我真的不明白如何明智地使用SPARK_CLASSPATH或SparkContext.addJar.我试过两个,看起来addJar不像以前那样工作.

在我的情况下,我试图在闭包或外面使用一些joda-time函数.如果我使用joda-time jar的路径设置SPARK_CLASSPATH,一切正常.但是,如果我删除SPARK_CLASSPATH并添加我的程序:

JavaSparkContext sc = new JavaSparkContext("spark://localhost:7077", "name", "path-to-spark-home", "path-to-the-job-jar");
sc.addJar("path-to-joda-jar");
Run Code Online (Sandbox Code Playgroud)

它不再起作用,虽然在日志中我可以看到:

14/03/17 15:32:57 INFO SparkContext: Added JAR /home/hduser/projects/joda-time-2.1.jar at http://127.0.0.1:46388/jars/joda-time-2.1.jar with timestamp 1395066777041
Run Code Online (Sandbox Code Playgroud)

并且立即:

Caused by: java.lang.NoClassDefFoundError: org/joda/time/DateTime
    at com.xxx.sparkjava1.SimpleApp.main(SimpleApp.java:57)
    ... 6 more
Caused by: java.lang.ClassNotFoundException: org.joda.time.DateTime
    at java.net.URLClassLoader$1.run(URLClassLoader.java:366)
Run Code Online (Sandbox Code Playgroud)

我以前认为SPARK_CLASSPATH正在为作业的驱动程序部分设置类路径,而SparkContext.addJar正在设置执行程序的类路径,但它似乎不再正确.

谁知道比我更好?

apache-spark

5
推荐指数
1
解决办法
1万
查看次数

纱线上的火花; 如何将指标发送到石墨槽?

我是新来的火花,我们在纱线上运行火花.我可以运行我的测试应用程序.我正在尝试收集Graphite中的spark指标.我知道要对metrics.properties文件进行哪些更改.但是我的spark应用程序将如何看待这个conf文件?

/xxx/spark/spark-0.9.0-incubating-bin-hadoop2/bin/spark-class org.apache.spark.deploy.yarn.Client --jar /xxx/spark/spark-0.9.0-incubating-bin-hadoop2/examples/target/scala-2.10/spark-examples_2.10-assembly-0.9.0-incubating.jar --addJars "hdfs://host:port/spark/lib/spark-assembly_2.10-0.9.0-incubating-hadoop2.2.0.jar" --class org.apache.spark.examples.Test --args yarn-standalone --num-workers 50 --master-memory 1024m --worker-memory 1024m --args "xx"
Run Code Online (Sandbox Code Playgroud)

我应该在哪里指定metrics.properties文件?

我对它做了这些修改:

*.sink.Graphite.class=org.apache.spark.metrics.sink.GraphiteSink
*.sink.Graphite.host=machine.domain.com
*.sink.Graphite.port=2003

master.source.jvm.class=org.apache.spark.metrics.source.JvmSource

worker.source.jvm.class=org.apache.spark.metrics.source.JvmSource

driver.source.jvm.class=org.apache.spark.metrics.source.JvmSource

executor.source.jvm.class=org.apache.spark.metrics.source.JvmSource
Run Code Online (Sandbox Code Playgroud)

hadoop scala apache-spark

5
推荐指数
2
解决办法
7890
查看次数

Elasticsearch 1.3. - 从Java调用自定义REST端点

我目前正在构建一个elasticsearch plugin公开REST端点(从这篇文章开始)

我可以这样调用我的端点curl:

curl -X POST 'http://my-es:9200/lt-dev_terminology_v1/english/_terminology?pretty=-d '{ 
   "segment": "database", 
   "analyzer": "en_analyzer"
}
Run Code Online (Sandbox Code Playgroud)

我的问题是如何java使用传输客户端调用相同的端点?你能指点我一些教程吗?

java rest elasticsearch

5
推荐指数
1
解决办法
652
查看次数

Elasticsearch:无法通过curl连接,奇怪的不一致行为

我正在运行Elasticsearch 1.3.4,通过Homebrew在Mac OS X 10.10上新安装:

$ brew install elasticsearch
$ elasticsearch
Run Code Online (Sandbox Code Playgroud)

http://localhost:9200/_cluster/state在浏览器中运行成功:

{
  "cluster_name": "elasticsearch_jbrukh",
  "version": 2,
  "master_node": "q6Jzcza_RwaVvc_1u95O1Q",
  "blocks": {},
  "nodes": {
    "q6Jzcza_RwaVvc_1u95O1Q": {
      "name": "Ethan Edwards",
      "transport_address": "inet[/127.0.0.1:9300]",
      "attributes": {}
    }
  },
  "metadata": {
    "templates": {},
    "indices": {}
  },
  "routing_table": {
    "indices": {}
  },
  "routing_nodes": {
    "unassigned": [],
    "nodes": {
      "q6Jzcza_RwaVvc_1u95O1Q": []
    }
  },
  "allocations": []
}
Run Code Online (Sandbox Code Playgroud)

但是,以下curl命令失败:

$ curl -XGET "http://localhost:9200/_cluster/state"
curl: (7) Failed to connect to localhost port 9200: Connection refused
Run Code Online (Sandbox Code Playgroud)

此外,curl命令间歇性地成功,但只有在从浏览器中点击该URL之后,它才会工作一次,然后再次因上述错误而再次失败.

我该如何解决这个问题?

curl elasticsearch

5
推荐指数
2
解决办法
2万
查看次数

以秒为单位转换 HH:mm:ss

我使用以下脚本从 yyyy-MM-dd HH:mm:ss 日期格式中提取 HH:mm:ss

import java.sql.Time

case class Transactions(creationTime: Time)

val formatter = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")

def parseTransac(line: String): (Transactions) = {
   val fields = line.split(',')
   val creationTime = new Time(formatter.parse(fields(0)).getTime())
   val transactions = Transactions(creationTime)
   (transactions)
}
Run Code Online (Sandbox Code Playgroud)

例如2009-01-15 15:45:23将返回15:45:23.

如何以秒为单位(56723)而不是 HH:mm:ss 获得结果

scala apache-spark

5
推荐指数
1
解决办法
3767
查看次数

是否有任何 pyspark 函数可以添加下个月,例如 DATE_ADD(date, Month(int type))

我是 Spark 的新手,是否有任何内置函数可以显示当前日期的下个月日期,例如今天是 2016 年 12 月 27 日,那么该函数将返回 2017 年 1 月 27 日。我已经使用了 date_add() 但没有添加月份的功能。我尝试过 date_add(date, 31) 但是如果这个月有 30 天怎么办?

spark.sql("select date_add(current_date(),31)") .show()
Run Code Online (Sandbox Code Playgroud)

谁能帮我解决这个问题。我需要为此编写自定义函数吗?因为我仍然没有找到任何内置代码提前感谢 Kalyan

python apache-spark apache-spark-sql pyspark

5
推荐指数
1
解决办法
1万
查看次数