我正在尝试导入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)
这不仅看起来很糟糕,我不想为每次测试都做到这一点."正确"的做法是什么?
我想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) 我需要在两个位置之间获得距离,但我需要像图中的蓝线一样得到距离.

我接下来尝试:
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)
但我得到红线距离.
我正在将一些机器学习算法(如线性回归,Logistic回归和朴素贝叶斯)应用于某些数据,但我试图避免使用RDD并开始使用DataFrame,因为RDD比pyspark下的Dataframe 慢(见图1).

我使用DataFrames的另一个原因是因为ml库有一个非常有用的类来调整模型,CrossValidator这个类在拟合之后返回一个模型,显然这个方法必须测试几个场景,然后返回一个拟合的模型(与参数的最佳组合).
我使用的集群不是那么大,数据相当大,有些适合需要几个小时,所以我想保存这些模型以便以后重用它们,但我还没有意识到,有什么我忽略的东西?
笔记:
如果有任何方法可以使用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代码,感觉有点迷失,我无法为自己的目的复制它们的功能.
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) 如果我有一个对象的实例,有没有办法检查我是否有一个单例对象而不是一个类的实例?有没有办法可以做到这一点?可能是一些反思API?我知道一个区别是单例对象的类名以a结尾$,但这不是一种严格的方法.
我想使用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中立即可读?
我试图用自定义阶段创建和保存管道.我需要使用一个添加column到我DataFrame的UDF.因此,我想知道是否有可能将一个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
我试图预测一个标签在每一行DataFrame,但不使用LinearRegressionModel的transform方法,由于醉翁之意不在酒,而不是我试图通过使用经典公式手动计算它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相同的特征 …
apache-spark ×8
pyspark ×5
scala ×5
python ×2
android ×1
dataframe ×1
google-maps ×1
hive ×1
rdd ×1