显示标签为“Spark”的博文。显示所有博文
显示标签为“Spark”的博文。显示所有博文

2018年4月25日星期三

spark history server 配置

参考 Spark入门 - History Server配置使用:http://callmesurprise.github.io/2016/11/13/Spark%E5%85%A5%E9%97%A8%20-%20history%20server/

同时需要把 spark.history.fs.cleaner.enabled 设置为 true,默认每天清理一次,最多保留七天的日志。参考:http://wxmimperio.tk/2016/01/22/Spark-JobHistory-Monitoring/

其他可以参考官方文档。

2018年1月29日星期一

spark中ALS算法的评测

看了下ALS的源码,如果设置了implicitPrefs,fit的过程优化的是原论文中 c_ui * (p_ui - x_u * y) ^ 2 + 正则项,计算出了 userFactor 和 itemFactor,也就是式中的x和y。recommend_for_all 或者 transform 来计算 prediction 时计算了 x*y 的预测值,所以预测的范围是0到1之间的regression。

所以如果设置了隐式反馈的参数为True,评测的时候就不能用rmse指标来判断了。

2018年1月18日星期四

pyspark将dataframe中的none转为空数组

from pyspark.sql.functions import *
from pyspark.sql.types import *

df = spark.createDataFrame([{'a': None, 'b': 2}, {'a': [2,3], 'b': 4}])
empty_array = udf(lambda :[], ArrayType(LongType()))

# solution 1
df.withColumn('a', coalesce(col('a'), empty_array())).show()

# solution 2
df.withColumn('a', when(isnull('a'), empty_array()).otherwise(df.a)).show()

2017年11月2日星期四

pyspark sql使用udf后yarn模式运行卡住

https://stackoverflow.com/questions/35157322/spark-dataframe-in-python-execution-stuck-when-using-udfs

我们用的pyspark 2.1.0版本,udf还有各种各样的问题,而且性能很差,只能转成rdd再做操作?

pyspark sql user defined function

参考https://docs.databricks.com/spark/latest/spark-sql/udf-in-python.html

举例:
>>> a = [{'a': 'a', 'b': 1}, {'a': 'aa', 'b': 2}]
>>> df = spark.createDataFrame(a)
[Row(a=u'a', b=1), Row(a=u'aa', b=2)]                                         
>>> def func(str):
...   return len(str) > 1
...
>>> from pyspark.sql.functions import udf
>>> from pyspark.sql.types import BooleanType
>>> func_udf = udf(func, BooleanType())
>>> df2 = df.filter(func_udf(df['a']))
>>> df2.collect()
[Row(a=u'aa', b=2)]

2017年10月31日星期二

spark合并小文件

spark使用FileUtil.copyMerge来进行小文件合并:https://hadoop.apache.org/docs/r2.7.1/api/org/apache/hadoop/fs/FileUtil.html

pyspark中dataframe union的一个问题

>>> a = [{'a': 1, 'b': 2}]
>>> x = spark.createDataFrame(a)
>>> b = sc.parallelize([(3, 4)])
>>> y = spark.createDataFrame(b, ['b', 'a'])
>>> x.collect()
[Row(a=1, b=2)]
>>> y.collect()
[Row(b=3, a=4)]
>>> z = x.union(y)
>>> z.collect()
[Row(a=1, b=2), Row(a=3, b=4)]

正确的结果应该是[Row(a=1, b=2), Row(a=4, b=3)],但实际输出的结果第二个Row的a和b反了。猜测DataFrame的union是按照顺序来的,并不是按照column的名称对应的。

Also as standard in SQL, this function resolves columns by position (not by name). Spark 2.3提供了unionByName可以解决问题,目前解决办法是把x与y的字段名排序要一样才行。

2017年10月30日星期一

spark读写mongodb的一个问题

从mongodb的某个collection中读取了df,做了一些操作后又overwrite写回该collection会有问题。因为在写的时候才action,猜测可能因为分布式的同时读写造成的问题。

问题确认:
将df cache,在回写之前先做一次action,让结果缓存到内存,然后再写mongo没有问题。

解决:
从一个collection读,写到另一个 collection

pyspark中判断DataFrame是否为空

if df.head() is not None:
  xxx

2017年10月27日星期五

spark 2.1.0 from_json使用中的问题

对于以下代码,spark2.2.0运行正常:
import json
from pyspark.sql import functions as f
from pyspark.sql.types import ArrayType, DoubleType, StringType, StructField, StructType
from pyspark.sql.functions import from_json

def func(value, score):
  values = {}
  for i in range(len(value)):
    if value[i] in values:
      values[value[i]] = values[value[i]] + score[i]
    else:
      values[value[i]] = score[i]
  res = []
  for k, v in values.items():
    res.append({'value': k, 'score': v})
  return json.dumps(res, ensure_ascii=False)

x = [{'user' : '86209203000295', 'domain' : 'music', 'subdomain' : 'artist', 'value' : 'xxx', 'score' : 0.8, 'ts' : '1508737410941'}, {'user' : '86209203000295', 'domain' : 'music', 'subdomain' : 'artist', 'value' : 'yyy', 'score' : 0.9, 'ts' : '1508737410941'}, {'user' : '86209203000685', 'domain' : 'music', 'subdomain' : 'artist', 'value' : 'zzz', 'score' : 0.8, 'ts' : '1508717416320'}]
df = spark.createDataFrame(x)
df = df.groupBy(df['user'], df['domain'], df['subdomain']).agg(f.collect_list(df['value']).alias('value'), f.collect_list(df['score']).alias('score'))
df = df.select(df['user'], df['domain'], df['subdomain'], f.UserDefinedFunction(func, StringType())(df['value'], df['score']).alias('values'))
df.collect()
schema = ArrayType(StructType([StructField('value', StringType()), StructField('score', DoubleType())]))
df = df.select(df['user'], df['domain'], df['subdomain'], from_json(df['values'], schema).alias('values'))
df.collect()

但是spark2.1.0运行报错:java.lang.ClassCastException: org.apache.spark.sql.types.ArrayType cannot be cast to org.apache.spark.sql.types.StructType

这个问题比较坑,2.1.0不支持ArrayType。

2017年9月13日星期三

获取spark-submit --files的文件


If you add your external files using "spark-submit --files" your files will be uploaded to this HDFS folder: hdfs://your-cluster/user/your-user/.sparkStaging/application_1449220589084_0508

application_1449220589084_0508 is an example of yarn application ID!

1. find the spark staging directory by below code: (but you need to have the hdfs uri and your username)

System.getenv("SPARK_YARN_STAGING_DIR"); --> .sparkStaging/application_1449220589084_0508

2. find the complete comma separated file paths by using:

System.getenv("SPARK_YARN_CACHE_FILES"); --> hdfs://yourcluster/user/hdfs/.sparkStaging/application_1449220589084_0508/spark-assembly-1.4.1.2.3.2.0-2950-hadoop2.7.1.2.3.2.0-2950.jar#__spark__.jar,hdfs://yourcluster/user/hdfs/.sparkStaging/application_1449220589084_0508/your-spark-job.jar#__app__.jar,hdfs://yourcluster/user/hdfs/.sparkStaging/application_1449220589084_0508/test_file.txt#test_file.txt


我的总结(以--files README.md为例):
方法1:按照上面所说,--files会把文件上传到hdfs的.sparkStagin/applicationId目录下,使用上面说的方法先获取到hdfs对应的这个目录,然后访问hdfs的这个文件。
spark.read().textFile(System.getenv("SPARK_YARN_STAGING_DIR") + "/README.md")解决。textFile不指定hdfs、file或者去其他前缀的话默认是hdfs://yourcluster/user/your_username下的相对路径。不知道是不是我使用的集群是这样设置的。

方法2:
SparkFiles.get(filePath),我获取的结果是:/hadoop/yarn/local/usercache/research/appcache/application_1504461219213_9796/spark-c39002ee-01a4-435f-8682-2ba5950de230/userFiles-e82a7f84-51b1-441a-a5e3-78bf3f4a8828/README.md,不知道为什么,无论本地还是hdfs都没有找到该文件。看了一下,本地是有/hadoop/yarn/local/usercache/research/...目录下的确有README.md。worker和driver的本地README.md路径不一样。
原因:
https://stackoverflow.com/questions/35865320/apache-spark-filenotfoundexception
https://stackoverflow.com/questions/41677897/how-to-get-path-to-the-uploaded-file
SparkFiles.get()获取的目录是driver node下的本地目录,所以sc.textFile无法在worker节点访问该目录文件。不能这么用。
"""I think that the main issue is that you are trying to read the file via the textFile method. What is inside the brackets of the textFile method is executed in the driver program. In the worker node only the code tobe run against an RDD is performed. When you type textFile what happens is that in your driver program it is created a RDD object with a trivial associated DAG.But nothing happens in the worker node."""


关于--files和addfile,可以看下这个问题:https://stackoverflow.com/questions/38879478/sparkcontext-addfile-vs-spark-submit-files

cluster模式下本地文件使用addFile是找不到文件的,因为只有本地有,所以必须使用--files上传。


结论:不要使用textFile读取--files或者addFile传来的文件。

SparkFiles.get出现NullPointerException错误

错误代码:
val serFile = SparkFiles.get("myobject.ser")

原因:SparkFiles.get只能在spark算子内使用:
sc.parallelize(1 to 100).map { i => SparkFiles.get("my.file") }.collect()

2017年3月2日星期四

传递java option的-D参数给spark-submit

参考http://stackoverflow.com/questions/28166667/how-to-pass-d-parameter-or-environment-variable-to-spark-job

我使用了com.typesafe.config,需要根据生产环境通过java option参数指定不同的config文件,saprk-submit增加如下选项:
--files your/config/file
--conf "spark.driver.extraJavaOptions=-Dconfig.resource=your_config_file.conf"
--conf "spark.executor.extraJavaOptions=-Dconfig.resource=your_config_file.conf"

在yarn-cluster模式下可行。
尝试了把--conf "spark.driver.extraJavaOptions" 换成了--driver-java-options,yarn-client模式依然出错。有时间再看看是什么问题。

2016年9月14日星期三

Spark中解析json

建议使用fastjson:
libraryDependencies += "com.alibaba" % "fastjson" % "1.2.24"


其他一些方案(不推荐):

scala.util.parsing.json.JSON:
解析:JSON.parseFull(x).get.asInstanceOf[Map[String, Any]]
生成:JSONObject(Map("field" -> value, ...))

或者com.fasterxml.jackson解析和生成json。

2016年9月13日星期二

Spark join的用法

用一个例子说明一下:
val a = sc.parallelize(List((1, "a"), (3, "c"), (5, "e"), (6, "f")))
val b = sc.parallelize(List((1, "a"), (2, "b"), (3, "c"), (4, "d")))

a.join(b).collect
    Array[(Int, (String, String))] = Array((1,(a,a)), (3,(c,c)))

a.leftOuterJoin(b).collect
    Array[(Int, (String, Option[String]))] = Array((1,(a,Some(a))), (6,(f,None)), (3,(c,Some(c))), (5,(e,None)))

a.rightOuterJoin(b).collect
    Array[(Int, (Option[String], String))] = Array((4,(None,d)), (1,(Some(a),a)), (3,(Some(c),c)), (2,(None,b)))

a.fullOuterJoin(b).collect
    Array[(Int, (Option[String], Option[String]))] = Array((4,(None,Some(d))), (1,(Some(a),Some(a))), (6,(Some(f),None)), (3,(Some(c),Some(c))), (5,(Some(e),None)), (2,(None,Some(b))))

a.subtractByKey(b).collect
    Array[(Int, String)] = Array((5,e), (6,f))

如果a或者b中有重复的key,join的结果会有多条。

2016年9月12日星期一

解决scala.reflect.api.JavaUniverse.runtimeMirror(Ljava/lang/ClassLoader;)Lscala/reflect/api/JavaMirrors$JavaMirror的问题

运行Spark的时候遇到这样一个问题:
java.lang.NoSuchMethodError: scala.reflect.api.JavaUniverse.runtimeMirror(Ljava/lang/ClassLoader;)Lscala/reflect/api/JavaMirrors$JavaMirror;

首先检查了一下我打包时候的Spark版本是基于scala 2.11的Spark 2.0.0,和spark-submit的版本是一致的,问题还是没有解决。最后发现是我的pom.xml中maven-scala-plugin的jvm版本设置的不对,修改了正确的jvm版本后运行成功。

2016年9月8日星期四

spark中dbscan使用中遇到的问题

Spark中目前没有集成dbscan聚类算法,见https://issues.apache.org/jira/browse/SPARK-5226

找了三个项目:
1、https://github.com/scalanlp/nak
nak中dbscan的输入必须是一个breeze.linalg.DenseMatrix的矩阵包含了所有数据,并不是我所需要的RDD[Vector],不可用。

2、https://github.com/irvingc/dbscan-on-spark
输入是一个RDD[Vector],但代码中只考虑了Vector的前两个值,所以该项目只能计算二维的dbscan,需要自己修改代码支持多维的Vector。

因为我用的Spark 2.0.0,而这个项目用的1.6.1,org.apache.spark.Logging类不存在了,我改为了org.apache.spark.internal.Logging。

运行时又遇到另外一个bug:
java.lang.NoSuchMethodError: scala.collection.immutable.$colon$colon.hd$1()Ljava/lang/Object;
看了一下代码,不知道为什么运行时不支持EvenSplitPartitioner.scala文件中的case (rectangle, count) :: rest => 这种List的写法,所以我将这段改为了获取List的head和tail,重新编译,提交spark运行,终于成功。

另外,该项目的repo中没有最新版本的jar,所以只能通过源码进行打包。

3、https://github.com/alitouka/spark_dbscan
目前只支持csv格式的文件作为输入,每行用逗号隔开。
同样,该项目是基于Spark 1.1.0写的,修改为2.0.0时要删去Logging类。

看了一下代码,如果不用IOHelper.readDataset方法,可以直接将数据转为RDD[Point]后进行计算。

运行了一下,发现vector维度为10000时候小数据量就报错,维度1000时大数据量也会出错,该问题目前还没有解决。

2016年7月29日星期五

Spark单元测试中使用mockito打桩

pom构建添加依赖:
    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-core_2.10</artifactId>
      <version>${spark.version}</version>
      <type>test-jar</type>
      <scope>test</scope>
    </dependency>
    <dependency>
      <groupId>org.mockito</groupId>
      <artifactId>mockito-all</artifactId>
      <version>1.9.5</version>
    </dependency>

sbt构建添加依赖:
"org.apache.spark" %% "spark-core" % "2.0.0" % "test" classifier "tests",
"org.mockito" % "mockito-all" % "1.9.5" % "test"


对class中的方法打桩(不能对class中变量打桩),使用spy。一个例子:
ReadFile.scala:
import org.apache.spark.SparkContext
import org.apache.spark.rdd.RDD

class ReadFile(sc: SparkContext) extends Serializable{
  def input: RDD[String] = {
    sc.textFile("hdfs://ip_address/xxx/data.txt")
  }

  def output = input.collect()
}

ReadFileSuite.scala:
import org.apache.spark.LocalSparkContext.withSpark
import org.apache.spark.{LocalSparkContext, SparkContext, SparkFunSuite}
import org.mockito.Mockito.{spy, when}

class ReadFileSuite extends SparkFunSuite with LocalSparkContext {
  test("test1") {
    withSpark(new SparkContext("local", "test")) { sc =>
      val data = spy(new ReadFile(sc))
      val stub_file = sc.textFile(getClass.getResource("/data.txt").getFile)
      when(data.input).thenReturn(stub_file)
      assert(data.output.length === 3)
    }
  }
}

ReadFileSuite中文件是从test/resources中读取的,而不是ReadFile中从hdfs读的路径。

如果class没有构造参数,可以使用mock(classOf[yourClass])创建mock的class。

参考了http://qiuguo0205.iteye.com/blog/1443344 和 http://qiuguo0205.iteye.com/blog/1456528