小编Alb*_*nto的帖子

在scala中导入spark.implicits._

我正在尝试导入spark.implicits._显然,这是scala中类中的一个对象.当我用这样的方法导入它时:

def f() = {
  val spark = SparkSession()....
  import spark.implicits._
}
Run Code Online (Sandbox Code Playgroud)

它工作正常,但我正在编写一个测试类,我想让这个导入可用于我尝试过的所有测试:

class SomeSpec extends FlatSpec with BeforeAndAfter {
  var spark:SparkSession = _

  //This won't compile
  import spark.implicits._

  before {
    spark = SparkSession()....
    //This won't either
    import spark.implicits._
  }

  "a test" should "run" in {
    //Even this won't compile (although it already looks bad here)
    import spark.implicits._

    //This was the only way i could make it work
    val spark = this.spark
    import spark.implicits._
  }
}
Run Code Online (Sandbox Code Playgroud)

这不仅看起来很糟糕,我不想为每次测试都做到这一点."正确"的做法是什么?

scala apache-spark

44
推荐指数
3
解决办法
3万
查看次数

使用列的长度过滤DataFrame

我想DataFrame使用与列长度相关的条件来过滤a ,这个问题可能很容易,但我没有在SO中找到任何相关的问题.

更具体的,我有一个DataFrame只有一个Column,其中ArrayType(StringType()),我要筛选的DataFrame使用长度filterer,我拍下面的一个片段.

df = sqlContext.read.parquet("letters.parquet")
df.show()

# The output will be 
# +------------+
# |      tokens|
# +------------+
# |[L, S, Y, S]|
# |[L, V, I, S]|
# |[I, A, N, A]|
# |[I, L, S, A]|
# |[E, N, N, Y]|
# |[E, I, M, A]|
# |[O, A, N, A]|
# |   [S, U, S]|
# +------------+

# But I want only the entries with length …
Run Code Online (Sandbox Code Playgroud)

python dataframe apache-spark apache-spark-sql pyspark

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

获取android中两个位置之间的距离?

我需要在两个位置之间获得距离,但我需要像图中的蓝线一样得到距离. picure

我接下来尝试:

public double getDistance(LatLng LatLng1, LatLng LatLng2) {
    double distance = 0;
    Location locationA = new Location("A");
    locationA.setLatitude(LatLng1.latitude);
    locationA.setLongitude(LatLng1.longitude);
    Location locationB = new Location("B");
    locationB.setLatitude(LatLng2.latitude);
    locationB.setLongitude(LatLng2.longitude);
    distance = locationA.distanceTo(locationB);

    return distance;
}
Run Code Online (Sandbox Code Playgroud)

但我得到红线距离.

android google-maps android-location

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

保存ML模型以备将来使用

我正在将一些机器学习算法(如线性回归,Logistic回归和朴素贝叶斯)应用于某些数据,但我试图避免使用RDD并开始使用DataFrame,因为RDD比pyspark下的Dataframe (见图1).

我使用DataFrames的另一个原因是因为ml库有一个非常有用的类来调整模型,CrossValidator这个类在拟合之后返回一个模型,显然这个方法必须测试几个场景,然后返回一个拟合的模型(与参数的最佳组合).

我使用的集群不是那么大,数据相当大,有些适合需要几个小时,所以我想保存这些模型以便以后重用它们,但我还没有意识到,有什么我忽略的东西?

笔记:

  • mllib的模型类有一个保存方法(即NaiveBayes),但mllib没有CrossValidator并使用RDD,所以我有预谋地避免它.
  • 目前的版本是spark 1.5.1.

apache-spark pyspark apache-spark-ml apache-spark-mllib

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

如何在Pyspark中使用Scala类

如果有任何方法可以使用Scala课程Pyspark,我一直在寻找一段时间,而且我没有找到任何关于这个主题的文档或指南.

假设我创建了一个简单的类,Scala它使用了一些库apache-spark,例如:

class SimpleClass(sqlContext: SQLContext, df: DataFrame, column: String) {
  def exe(): DataFrame = {
    import sqlContext.implicits._

    df.select(col(column))
  }
}
Run Code Online (Sandbox Code Playgroud)
  • 有没有可能的方法来使用这个类Pyspark
  • 太难了吗?
  • 我必须创建一个.py文件吗?
  • 是否有任何指南说明如何做到这一点?

顺便说一句,我也查看了spark代码,感觉有点迷失,我无法为自己的目的复制它们的功能.

python scala apache-spark apache-spark-sql pyspark

19
推荐指数
2
解决办法
8355
查看次数

Apache Spark:使用RDD.aggregateByKey()的RDD.groupByKey()的等效实现是什么?

Apache Spark pyspark.RDDAPI文档提到groupByKey()效率低下.相反,它是推荐使用reduceByKey(),aggregateByKey(),combineByKey(),或foldByKey()代替.这将导致在shuffle之前在worker中进行一些聚合,从而减少跨工作人员的数据混乱.

给定以下数据集和groupByKey()表达式,什么是等效且有效的实现(减少的跨工作者数据混洗),它不使用groupByKey(),但提供相同的结果?

dataset = [("a", 7), ("b", 3), ("a", 8)]
rdd = (sc.parallelize(dataset)
       .groupByKey())
print sorted(rdd.mapValues(list).collect())
Run Code Online (Sandbox Code Playgroud)

输出:

[('a', [7, 8]), ('b', [3])]
Run Code Online (Sandbox Code Playgroud)

apache-spark rdd pyspark

11
推荐指数
1
解决办法
8486
查看次数

在scala中,有没有办法检查实例是否是单例对象?

如果我有一个对象的实例,有没有办法检查我是否有一个单例对象而不是一个类的实例?有没有办法可以做到这一点?可能是一些反思API?我知道一个区别是单例对象的类名以a结尾$,但这不是一种严格的方法.

scala

10
推荐指数
1
解决办法
808
查看次数

将Spark数据框保存到Hive:table不可读,因为"镶木地板不是SequenceFile"

我想使用PySpark将Spark(v 1.3.0)数据框中的数据保存到Hive表中.

文件规定:

"spark.sql.hive.convertMetastoreParquet:当设置为false时,Spark SQL将使用Hive SerDe作为镶木桌而不是内置支持."

看看Spark教程,似乎可以设置这个属性:

from pyspark.sql import HiveContext

sqlContext = HiveContext(sc)
sqlContext.sql("SET spark.sql.hive.convertMetastoreParquet=false")

# code to create dataframe

my_dataframe.saveAsTable("my_dataframe")
Run Code Online (Sandbox Code Playgroud)

但是,当我尝试查询Hive中保存的表时,它返回:

hive> select * from my_dataframe;
OK
Failed with exception java.io.IOException:java.io.IOException: 
hdfs://hadoop01.woolford.io:8020/user/hive/warehouse/my_dataframe/part-r-00001.parquet
not a SequenceFile
Run Code Online (Sandbox Code Playgroud)

如何保存表格,使其在Hive中立即可读?

hive apache-spark apache-spark-sql pyspark

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

如何从UDF创建自定义Transformer?

我试图用自定义阶段创建和保存管道.我需要使用一个添加column到我DataFrameUDF.因此,我想知道是否有可能将一个UDF或类似的动作转换成一个Transformer

我的自定义UDF看起来像这样,我想学习如何使用UDF自定义Transformer.

def getFeatures(n: String) = {
    val NUMBER_FEATURES = 4  
    val name = n.split(" +")(0).toLowerCase
    ((1 to NUMBER_FEATURES)
         .filter(size => size <= name.length)
         .map(size => name.substring(name.length - size)))
} 

val tokenizeUDF = sqlContext.udf.register("tokenize", (name: String) => getFeatures(name))
Run Code Online (Sandbox Code Playgroud)

scala user-defined-functions apache-spark apache-spark-sql apache-spark-ml

9
推荐指数
1
解决办法
4217
查看次数

LogisticRegressionModel手动预测

我试图预测一个标签在每一行DataFrame,但不使用LinearRegressionModeltransform方法,由于醉翁之意不在酒,而不是我试图通过使用经典公式手动计算它1 / (1 + e^(-h?(x))),请注意,我是从复制的代码Apache Spark的存储库并将几乎所有东西从private对象复制BLAS到它的公共版本中.PD:我没有使用任何regParam,我只是安装了模型.

//Notice that I had to obtain intercept, and coefficients from my model
val intercept = model.intercept
val coefficients = model.coefficients

val margin: Vector => Double = (features) => {
  BLAS.dot(features, coefficients) + intercept
}

val score: Vector => Double = (features) => {
  val m = margin(features)
  1.0 / (1.0 + math.exp(-m))
}
Run Code Online (Sandbox Code Playgroud)

在定义了这些函数并获得模型的参数后,我创建了一个UDF来计算预测(它接收与a相同的特征 …

scala logistic-regression apache-spark

9
推荐指数
1
解决办法
340
查看次数