我想使用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中立即可读?
这是问题(链接:http://opc.iarcs.org.in/index.php/problems/FINDPERM):
数字1,...,N的排列是这些数字的重新排列.例如,
2 4 5 1 7 6 3 8
是1,2,...,8的置换.当然,
1 2 3 4 5 6 7 8
也是1,2,...,8的置换.
与N的每个排列相关联的是长度为N的正整数的特殊序列,称为其反转序列.该序列的第i个元素是严格小于i的数字j的数量,并且在该排列中出现在i的右侧.对于置换
2 4 5 1 7 6 3 8
,反转序列是
0 1 0 2 2 1 2 0
第二个元素是1,因为1严格小于2,在这个排列中它出现在2的右边.类似地,第5个元素是2,因为1和3严格小于5但在此排列中出现在5的右侧,依此类推.
作为另一个例子,置换的反转序列
8 7 6 5 4 3 2 1
是
0 1 2 3 4 5 6 7
在这个问题中,你将得到一些置换的反演序列.您的任务是从此序列重建置换.
我想出了这段代码:
#include <iostream>
using namespace std;
void insert(int key, int *array, int value , int size){
int i …Run Code Online (Sandbox Code Playgroud) 我有数千万行数据.是否可以使用火花流在一周或一天内分析所有这些?在数据量方面,火花流的限制是什么?我不确定什么是上限,什么时候我应该将它们放入我的数据库,因为Stream可能无法再处理它们了.我也有不同的时间窗口1,3,6小时等,我使用窗口操作来分隔数据.
请在下面找到我的代码:
conf = SparkConf().setAppName(appname)
sc = SparkContext(conf=conf)
ssc = StreamingContext(sc,300)
sqlContext = SQLContext(sc)
channels = sc.cassandraTable("abc","channels")
topic = 'abc.crawled_articles'
kafkaParams = {"metadata.broker.list": "0.0.0.0:9092"}
category = 'abc.crawled_article'
category_stream = KafkaUtils.createDirectStream(ssc, [category], kafkaParams)
category_join_stream = category_stream.map(lambda x:read_json(x[1])).filter(lambda x:x!=0).map(lambda x:categoryTransform(x)).filter(lambda x:x!=0).map(lambda x:(x['id'],x))
article_stream = KafkaUtils.createDirectStream(ssc, [topic], kafkaParams)
article_join_stream=article_stream.map(lambda x:read_json(x[1])).filter(lambda x: x!=0).map(lambda x:TransformInData(x)).filter(lambda x: x!=0).flatMap(lambda x:(a for a in x)).map(lambda x:(x['id'].encode("utf-8") ,x))
#axes topic integration the article and the axes
axes_topic = 'abc.crawled_axes'
axes_stream = KafkaUtils.createDirectStream(ssc, [axes_topic], kafkaParams)
axes_join_stream = axes_stream.filter(lambda x:'delete' not …Run Code Online (Sandbox Code Playgroud) 有没有漂亮的代码可以在 Pyspark 中执行不区分大小写的连接?就像是:
df3 = df1.join(df2,
["col1", "col2", "col3"],
"left_outer",
"case-insensitive")
Run Code Online (Sandbox Code Playgroud)
或者你对此的工作解决方案是什么?
是否可以使用Shapely将 aMultipolygon转换为Polygon填充所有孔或缺失内部区域的 a ?我已经尝试了一段时间,但在文档中找不到它。下图显示了一个多边形示例,其中包含我要填充的孔和要删除的正方形。
我在想,如果有可能创建一个UDF接收两个参数的Column和另一个变量(Object,Dictionary或任何其他类型),然后做一些操作,并返回结果.
实际上,我试图这样做,但我有一个例外.因此,我想知道是否有办法避免这个问题.
df = sqlContext.createDataFrame([("Bonsanto", 20, 2000.00),
("Hayek", 60, 3000.00),
("Mises", 60, 1000.0)],
["name", "age", "balance"])
comparatorUDF = udf(lambda c, n: c == n, BooleanType())
df.where(comparatorUDF(col("name"), "Bonsanto")).show()
Run Code Online (Sandbox Code Playgroud)
我收到以下错误:
AnalysisException:u"无法解析'Bonsanto'给定的输入列名称,年龄,余额;"
所以很明显,UDF"看到" string"Bonsanto"作为列名,实际上我正在尝试将记录值与第二个参数进行比较.
另一方面,我知道可以在where子句中使用一些运算符(但实际上我想知道它是否可以使用UDF),如下所示:
df.where(col("name") == "Bonsanto").show()
#+--------+---+-------+
#| name|age|balance|
#+--------+---+-------+
#|Bonsanto| 20| 2000.0|
#+--------+---+-------+
Run Code Online (Sandbox Code Playgroud) python user-defined-functions apache-spark apache-spark-sql pyspark
带有独立主节点和2个工作节点的Spark独立群集,每个工作节点上有4个cpu核心。所有工人共有8个核心。
通过spark-submit运行以下命令时(未设置spark.default.parallelism)
val myRDD = sc.parallelize(1 to 100000)
println("Partititon size - " + myRDD.partitions.size)
val totl = myRDD.reduce((x, y) => x + y)
println("Sum - " + totl)
Run Code Online (Sandbox Code Playgroud)
返回分区大小的值2。
通过连接到Spark独立集群使用spark-shell时,相同代码返回正确的分区大小8。
可能是什么原因?
谢谢。
我试图使用Spark来处理简单的图形问题.我在Spark源文件夹中找到了一个示例程序:transitive_closure.py,它在一个图形中计算传递闭包,不超过200个边和顶点.但是在我自己的笔记本电脑中,它运行超过10分钟并且不会终止.我使用的命令行是:spark-submit transitive_closure.py.
我想知道为什么即使计算这么小的传递闭包结果,火花也是如此之慢?这是常见的情况吗?有没有我想念的配置?
该程序如下所示,可以在他们网站的spark install文件夹中找到.
from __future__ import print_function
import sys
from random import Random
from pyspark import SparkContext
numEdges = 200
numVertices = 100
rand = Random(42)
def generateGraph():
edges = set()
while len(edges) < numEdges:
src = rand.randrange(0, numEdges)
dst = rand.randrange(0, numEdges)
if src != dst:
edges.add((src, dst))
return edges
if __name__ == "__main__":
"""
Usage: transitive_closure [partitions]
"""
sc = SparkContext(appName="PythonTransitiveClosure")
partitions = int(sys.argv[1]) if len(sys.argv) > 1 else 2
tc = sc.parallelize(generateGraph(), partitions).cache()
# Linear …Run Code Online (Sandbox Code Playgroud) 我正在使用MLlib的矩阵分解来向用户推荐项目.我有一个很大的隐式交互矩阵,M = 2000万用户和N = 50k项目.在训练模型之后,我想获得每个用户的推荐列表(例如200).我尝试过recommendProductsForUsers,MatrixFactorizationModel但它非常慢(跑了9个小时,但距离完成还很远.我正在测试50个执行器,每个都有8g内存).这可能是预期的,因为recommendProductsForUsers需要计算所有M*N用户项目交互并获得每个用户的顶部.
我会尝试使用更多的执行程序,但是从我在Spark UI上的应用程序细节中看到的,我怀疑它可以在几小时或一天完成,即使我有1000个执行程序(9小时之后它仍然在flatmap这里https:// github. com/apache/spark/blob/master/mllib/src/main/scala/org/apache/spark/mllib/recommendation/MatrixFactorizationModel.scala#L279-L289,10000个总任务,只有~200个完成)还有其他的东西吗除了增加执行者数量之外,我可以调整以加快推荐过程吗?
这是示例代码:
val data = input.map(r => Rating(r.getString(0).toInt, r.getString(1).toInt, r.getLong(2))).cache
val rank = 20
val alpha = 40
val maxIter = 10
val lambda = 0.05
val checkpointIterval = 5
val als = new ALS()
.setImplicitPrefs(true)
.setCheckpointInterval(checkpointIterval)
.setRank(rank)
.setAlpha(alpha)
.setIterations(maxIter)
.setLambda(lambda)
val model = als.run(ratings)
val recommendations = model.recommendProductsForUsers(200)
recommendations.saveAsTextFile(outdir)
Run Code Online (Sandbox Code Playgroud) scala collaborative-filtering apache-spark apache-spark-mllib
在进行了一些转换之后dataframe,我得到了以下内容,在这种情况下如何通过列获取前n条记录short_name并使用其他作为指标frequency.我读过这篇文章,但两个解决方案的问题是他们摆脱了列product_name,他们只保留了分组列,我需要保留它们.
short_name product_id frequency
Yoghurt y cereales 975009684 32
Yoghurt y cereales 975009685 21
Yoghurt y cereales 975009700 16
Yoghurt y Cereales 21097 16
Yoghurt Bebible 21329 68
Yoghurt Bebible 21328 67
Yoghurt Bebible 21500 31
Run Code Online (Sandbox Code Playgroud)