小编mar*_*tin的帖子

PySpark按条件计算值

我有一个DataFrame,这里有一个片段:

[['u1', 1], ['u2', 0]]
Run Code Online (Sandbox Code Playgroud)

基本上是一个名为的字符串字段f,对于第二个元素(is_fav)为1或0 .

我需要做的是分组第一个字段并计算1和0的出现次数.我希望做类似的事情

num_fav = count((col("is_fav") == 1)).alias("num_fav")

num_nonfav = count((col("is_fav") == 0)).alias("num_nonfav")

df.groupBy("f").agg(num_fav, num_nonfav)
Run Code Online (Sandbox Code Playgroud)

它不能正常工作,我在两种情况下都得到相同的结果,这相当于组中项目的计数,因此过滤器(无论是1还是0)似乎被忽略.这取决于count工作原理吗?

python apache-spark pyspark

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

PySpark:在写入而不是多个部分文件时吐出单个文件

有没有办法阻止PySpark在将DataFrame写入JSON文件时创建几个小文件?

如果我跑:

 df.write.format('json').save('myfile.json')
Run Code Online (Sandbox Code Playgroud)

要么

df1.write.json('myfile.json')
Run Code Online (Sandbox Code Playgroud)

它创建了名为的文件夹myfile,在其中我找到了几个名为part-***HDFS的小文件.是否有可能让它吐出一个文件而不是?

python amazon-s3 apache-spark apache-spark-sql pyspark

7
推荐指数
1
解决办法
9603
查看次数

Spark如何决定如何划分RDD?

假设我创建了这样一个RDD(我正在使用Pyspark):

list_rdd = sc.parallelize(xrange(0, 20, 2), 6)
Run Code Online (Sandbox Code Playgroud)

然后我用glom()方法打印分区元素并获取

[[0], [2, 4], [6, 8], [10], [12, 14], [16, 18]]
Run Code Online (Sandbox Code Playgroud)

Spark如何决定如何对我的列表进行分区?这些元素的具体选择来自哪里?它可以以不同的方式耦合它们,留下除0和10之外的其他元素,以创建6个请求的分区.在第二次运行时,分区是相同的.

使用更大的范围,29个元素,我得到2个元素的模式分区,然后是三个元素:

list_rdd = sc.parallelize(xrange(0, 30, 2), 6)
[[0, 2], [4, 6, 8], [10, 12], [14, 16, 18], [20, 22], [24, 26, 28]]
Run Code Online (Sandbox Code Playgroud)

使用较小范围的9个元素

list_rdd = sc.parallelize(xrange(0, 10, 2), 6)
[[], [0], [2], [4], [6], [8]]
Run Code Online (Sandbox Code Playgroud)

所以我推断Spark是通过将列表拆分为一个配置来生成分区,其中最小可能后跟更大的集合,并重复.

问题是这个选择背后是否有原因,这是非常优雅的,但它是否也提供了性能优势?

apache-spark rdd pyspark

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

PySpark:当列是列表时,将列添加到DataFrame

我读过类似的问题,但是找不到我特定问题的解决方案。

我有一个清单

l = [1, 2, 3]
Run Code Online (Sandbox Code Playgroud)

和一个DataFrame

df = sc.parallelize([
    ['p1', 'a'],
    ['p2', 'b'],
    ['p3', 'c'],
]).toDF(('product', 'name'))
Run Code Online (Sandbox Code Playgroud)

我想获得一个新的DataFrame,其中将列表l添加为另一列,即

+-------+----+---------+
|product|name| new_col |
+-------+----+---------+
|     p1|   a|     1   |
|     p2|   b|     2   |
|     p3|   c|     3   |
+-------+----+---------+
Run Code Online (Sandbox Code Playgroud)

与JOIN接触的方法,我当时在DF与

 sc.parallelize([[1], [2], [3]])
Run Code Online (Sandbox Code Playgroud)

失败了。使用中的方法withColumn,如

new_df = df.withColumn('new_col', l)
Run Code Online (Sandbox Code Playgroud)

由于列表不是Column对象而失败。

python dataframe pyspark

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

在嵌套字段上加入PySpark DataFrames

我想在这两个PySpark DataFrame之间执行连接:

from pyspark import SparkContext
from pyspark.sql.functions import col

sc = SparkContext()

df1 = sc.parallelize([
    ['owner1', 'obj1', 0.5],
    ['owner1', 'obj1', 0.2],
    ['owner2', 'obj2', 0.1]
]).toDF(('owner', 'object', 'score'))

df2 = sc.parallelize(
          [Row(owner=u'owner1',
           objects=[Row(name=u'obj1', value=Row(fav=True, ratio=0.3))])]).toDF()
Run Code Online (Sandbox Code Playgroud)

在已经加入到在对象,即字段的名称进行命名的内部对象为DF2和对象为DF1.

我可以在嵌套字段上执行SELECT,如

df2.where(df2.owner == 'owner1').select(col("objects.value.ratio")).show()
Run Code Online (Sandbox Code Playgroud)

但是我无法运行此连接:

df2.alias('u').join(df1.alias('s'), col('u.objects.name') == col('s.object'))
Run Code Online (Sandbox Code Playgroud)

返回错误

pyspark.sql.utils.AnalysisException:由于数据类型不匹配,u"无法解析'(objects.name = cast(object as double))''''(objects.name = cast(object as double))'中的不同类型(数组和双).;"

任何想法如何解决这个问题?

join dataframe apache-spark apache-spark-sql pyspark

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

Django模型 - 更改选项字段中的选项并进行迁移

我需要在模型状态字段中更改选项的名称,这需要成为

   STATUS = Choices(
    ('option_A', 'Option A'),
    ('option_B', 'Option B'),
   )
Run Code Online (Sandbox Code Playgroud)

在此更改之前,我有相同的选项,但名称不同.现在我改变了项目中的所有内容以尊重新名称,但问题是更新数据库以显示这些新名称.我使用South进行数据迁移,据我所知,如果你需要在数据库中添加或删除一个列,那么让它编写一个自动迁移是相当简单的,但我找不到一种方法来对我进行更新.现有专栏.

我使用Django 1.6.

django django-models django-south django-migrations

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

Numpy中的直方图日期时间对象

我有一个datetime对象数组,我想用Python直方图.

Numpy 直方图方法不接受日期时间,抛出的错误是

File "/usr/lib/python2.7/dist-packages/numpy/lib/function_base.py", line 176, in histogram
mn, mx = [mi+0.0 for mi in range]
TypeError: unsupported operand type(s) for +: 'datetime.datetime' and 'float'
Run Code Online (Sandbox Code Playgroud)

除了手动转换datetime对象之外,还有其他方法可以执行此操作吗?

python datetime numpy histogram

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

通过在Jupyter中内联重置Jupyter中的Matplotlib图形大小

这个问题更像是一种好奇心.

要在matplotlib中将默认的图形大小更改为自定义图形,可以使用

from matplotlib import rcParams
from matplotlib import pyplot as plt
rcParams['figure.figsize'] = 15, 9
Run Code Online (Sandbox Code Playgroud)

之后,图形会显示所选尺寸.

现在,我发现了一些新的东西(从未发生过/现在才注意到):在Jupyter笔记本中,将matplotlib内联为

%matplotlib inline
Run Code Online (Sandbox Code Playgroud)

这显然会覆盖rcParams字典,恢复图形大小的默认值.因此,为了能够设置大小,我必须matplotlib 在更改rcParams字典的值之前内联.

我使用的是Mac OS 10.11.6,matplotlib 1.5.1版,Python 2.7.10,Jupyter 4.1.

python plot matplotlib jupyter

5
推荐指数
1
解决办法
1286
查看次数

配置Ipython的后端以使用带代码的视网膜显示模式

我正在使用代码配置Jupyter笔记本,因为我有一个包含大量笔记本的repo,并希望保持所有样式的一致性,而不必在每个开头都写出冗长的设置.这样,我所做的就是有一个配置CSS的方法,一个用于设置Matplotlib,另一个用于配置Ipython.

我按照这种方式配置我的笔记本而不是依赖于文档的配置文件的原因有两个:

  1. 我公开分享这个笔记本电脑的回购,我希望我的所有配置都可见
  2. 我想保留这些配置仅仅是我正在创建的这个回购

作为示例,设置CSS的方法看起来像

def set_css_style(css_file_path='../styles_files/custom.css'):

    styles = open(css_file_path, "r").read()
    return HTML(styles)
Run Code Online (Sandbox Code Playgroud)

我在每个笔记本的开头用它来调用它set_css_style().同样,我有这个方法来配置Ipython的细节:

def config_ipython():

    InteractiveShell.ast_node_interactivity = "all"
Run Code Online (Sandbox Code Playgroud)

以上都使用进口

from IPython.core.display import HTML
from IPython.core.interactiveshell import InteractiveShell
Run Code Online (Sandbox Code Playgroud)

目前,可以看出,配置Ipython的方法只包含指令,以便当我在单元格中的多行中键入变量名称时,我不需要添加一个print以使它们全部被打印.

我的问题是如何转换Jupyter魔术命令以获得数字代码的视网膜显示质量.这样的命令是

%config InlineBackend.figure_format = 'retina'
Run Code Online (Sandbox Code Playgroud)

从Ipython的文档中我找不到如何在一个方法中调用这个指令,即找不到InlineBackend生活的地方.

我只想将此配置行添加到config_ipython上面的方法中,是否可能?

ipython retina-display ipython-magic jupyter-notebook

5
推荐指数
1
解决办法
4880
查看次数

Boto3 逐行从 S3 键读取文件内容

使用 boto3,您可以从 S3 中的某个位置读取文件内容,根据存储桶名称和密钥(这假设是初步的import boto3)

s3 = boto3.resource('s3')

content = s3.Object(BUCKET_NAME, S3_KEY).get()['Body'].read()
Run Code Online (Sandbox Code Playgroud)

这将返回一个字符串类型。我需要获取的特定文件恰好是一组类似字典的对象,每行一个。所以它不是 JSON 格式。我不想将其作为字符串读取,而是将其作为文件对象流式传输并逐行读取;除了首先在本地下载文件之外,找不到其他方法来执行此操作

s3 = boto3.resource('s3')

bucket = s3.Bucket(BUCKET_NAME)

filename = 'my-file'
bucket.download_file(S3_KEY, filename)

f = open('my-file')
Run Code Online (Sandbox Code Playgroud)

我要问的是是否可以对文件进行这种类型的控制,而不必先在本地下载它?

python amazon-s3 amazon-web-services boto3

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