小编Alb*_*nto的帖子

将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万
查看次数

我需要一个更好的算法来解决这个问题

这是问题(链接: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)

c++ algorithm sequences

8
推荐指数
1
解决办法
1280
查看次数

在数据量方面,火花流的限制是什么?

我有数千万行数据.是否可以使用火花流在一周或一天内分析所有这些?在数据量方面,火花流的限制是什么?我不确定什么是上限,什么时候我应该将它们放入我的数据库,因为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)

datastax-enterprise apache-spark spark-streaming

8
推荐指数
1
解决办法
997
查看次数

Pyspark - 如何进行不区分大小写的数据帧连接?

有没有漂亮的代码可以在 Pyspark 中执行不区分大小写的连接?就像是:

df3 = df1.join(df2, 
               ["col1", "col2", "col3"],
               "left_outer",
               "case-insensitive")
Run Code Online (Sandbox Code Playgroud)

或者你对此的工作解决方案是什么?

apache-spark apache-spark-sql pyspark

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

在 Python 中将多多边形转换为多边形

是否可以使用Shapely将 aMultipolygon转换为Polygon填充所有孔或缺失内部区域的 a ?我已经尝试了一段时间,但在文档中找不到它。下图显示了一个多边形示例,其中包含我要填充的孔和要删除的正方形。

带孔的多边形

python polygons shapely

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

如何将常量值传递给Python UDF?

我在想,如果有可能创建一个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

6
推荐指数
1
解决办法
5666
查看次数

用于火花提交的Parallelize RDD的spark.default.parallelism默认为2

带有独立主节点和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。

可能是什么原因?

谢谢。

scala apache-spark

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

Spark示例程序运行速度很慢

我试图使用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)

performance transitive-closure apache-spark pyspark

6
推荐指数
1
解决办法
1927
查看次数

加快Spark MLLib中大型数据集的协同过滤

我正在使用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

6
推荐指数
1
解决办法
761
查看次数

在保留所有列的pandas中获取每个类别的前n个值

在进行了一些转换之后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)

python pandas

6
推荐指数
3
解决办法
1989
查看次数