在具有多个工作节点的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?
我无法理解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) 我想在数据框中添加一个新列,其值由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函数为列生成随机值?
使用动态资源分配:
当执行程序空闲超过 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?对数据库的阻塞调用算作空闲吗?
我正在使用一个独立的火花群,一个主人和两个工人.我真的不明白如何明智地使用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正在设置执行程序的类路径,但它似乎不再正确.
谁知道比我更好?
我是新来的火花,我们在纱线上运行火花.我可以运行我的测试应用程序.我正在尝试收集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) 我目前正在构建一个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使用传输客户端调用相同的端点?你能指点我一些教程吗?
我正在运行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之后,它才会工作一次,然后再次因上述错误而再次失败.
我该如何解决这个问题?
我使用以下脚本从 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 获得结果
我是 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