- android - 多次调用 OnPrimaryClipChangedListener
- android - 无法更新 RecyclerView 中的 TextView 字段
- android.database.CursorIndexOutOfBoundsException : Index 0 requested, 光标大小为 0
- android - 使用 AppCompat 时,我们是否需要明确指定其 UI 组件(Spinner、EditText)颜色
我一直在使用 PySpark 在相当大的 RDD 中运行大量计算,其每个 block 如下所示:
ID CHK C1 Flag1 V1 V2 C2 Flag2 V3 V4
341 10 100 TRUE 10 10 150 FALSE 10 14
341 9 100 TRUE 10 10 150 FALSE 10 14
341 8 100 TRUE 14 14 150 FALSE 10 14
341 7 100 TRUE 14 14 150 FALSE 10 14
341 6 100 TRUE 14 14 150 FALSE 10 14
341 5 100 TRUE 14 14 150 FALSE 10 14
341 4 100 TRUE 14 14 150 FALSE 12 14
341 3 100 TRUE 14 14 150 FALSE 14 14
341 2 100 TRUE 14 14 150 FALSE 14 14
341 1 100 TRUE 14 14 150 FALSE 14 14
341 0 100 TRUE 14 14 150 FALSE 14 14
我有很多 ID(它取决于 C1 值,例如从 100 到 130 等等,对于许多 C1,对于每个整数,我有一组 11 行,就像上面的那样)并且我有很多 ID .我需要做的是在每一行的组中应用一个公式并添加两列来计算:
D1 = ((row.V1 - prev_row.V1)/2)/((row.V2 + prev_row.V2)/2)
D2 = ((row.V3 - prev_row.V3)/2)/((row.V4 + prev_row.V4)/2)
我所做的(正如我在这篇有用的文章中发现的:https://arundhaj.com/blog/calculate-difference-with-previous-row-in-pyspark.html)是定义一个窗口:
my_window = Window.partitionBy().orderBy(desc("CHK"))
我为每个中间计算创建了一个“临时”列:
df = df.withColumn("prev_V1", lag(df.V1).over(my_window))
df = df.withColumn("prev_V21", lag(df.TA1).over(my_window))
df = df.withColumn("prev_V3", lag(df.SSQ2).over(my_window))
df = df.withColumn("prev_V4", lag(df.TA2).over(my_window))
df = df.withColumn("Sub_V1", F.when(F.isnull(df.V1 - df.prev_V1), 0).otherwise((df.V1 - df.prev_V1)/2))
df = df.withColumn("Sub_V2", (df.V2 + df.prev_V2)/2)
df = df.withColumn("Sub_V3", F.when(F.isnull(df.V3 - df.prev_V3), 0).otherwise((df.V3 - df.prev_V3)/2))
df = df.withColumn("Sub_V4", (df.V4 + df.prev_V4)/2)
df = df.withColumn("D1", F.when(F.isnull(df.Sub_V1 / df.Sub_V2), 0).otherwise(df.Sub_V1 / df.Sub_V2))
df = df.withColumn("D2", F.when(F.isnull(df.Sub_V3 / df.Sub_V4), 0).otherwise(df.Sub_V3 / df.Sub_V4))
最后我摆脱了临时列:
final_df = df.select(*columns_needed)
花了很长时间,我不断得到:
WARN WindowExec: No Partition Defined for Window operation! Moving all data to a single partition, this can cause serious performance degradation.
我知道我没有正确地执行此操作,因为上面的代码块位于几个 for 循环内,以便对所有 ID 进行计算,即循环使用:
unique_IDs = list(df1.toPandas()['ID'].unique())
但在深入了解 PySpark 窗口函数后,我相信通过正确设置窗口 partitionBy(),我可以更轻松地获得相同的结果。
我查看了 Avoid performance impact of a single partition mode in Spark window functions但我仍然不确定如何正确设置我的窗口分区以使其正常工作。
有人可以就如何解决这个问题提供一些帮助或见解吗?
谢谢
最佳答案
我假设公式必须应用于每组 ID(这就是我选择按“ID”进行分区的原因)。
你可以避免像这样使用'temp'列:
# used to define the lag of a specific column
w_lag=Window.partitionBy("id","C1").orderBy(desc('chk'))
df = df.withColumn('D1',((df.V1-F.lag(df.V1).over(w_lag))/2)\
/((df.V2+F.lag(df.V2).over(w_lag))/2))
df = df.withColumn('D2',((df.V3-F.lag(df.V3).over(w_lag))/2)\
/((df.V4+F.lag(df.V4).over(w_lag))/2))
结果是:
+---+---+---+-----+---+---+---+-----+---+---+-------------------+-------------------+
| id|chk| C1|Flag1| V1| V2| C2|Flag2| V3| V4| D1| D2|
+---+---+---+-----+---+---+---+-----+---+---+-------------------+-------------------+
|341| 10|100| true| 10| 10|150| true| 10| 14| null| null|
|341| 9|100| true| 10| 10|150| true| 10| 14| 0.0| 0.0|
|341| 8|100| true| 14| 14|150| true| 10| 14|0.16666666666666666| 0.0|
|341| 7|100| true| 14| 14|150| true| 10| 14| 0.0| 0.0|
|341| 6|100| true| 14| 14|150| true| 10| 14| 0.0| 0.0|
|341| 5|100| true| 14| 14|150| true| 10| 14| 0.0| 0.0|
|341| 4|100| true| 14| 14|150| true| 12| 14| 0.0|0.07142857142857142|
|341| 3|100| true| 14| 14|150| true| 14| 14| 0.0|0.07142857142857142|
|341| 2|100| true| 14| 14|150| true| 14| 14| 0.0| 0.0|
|341| 1|100| true| 14| 14|150| true| 14| 14| 0.0| 0.0|
|341| 0|100| true| 14| 14|150| true| 14| 14| 0.0| 0.0|
+---+---+---+-----+---+---+---+-----+---+---+-------------------+-------------------+
关于python - PySpark 窗口函数理解,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/46546634/
https://github.com/mattdiamond/Recorderjs/blob/master/recorder.js中的代码 我不明白 JavaScript 语法,比如 (functio
在 iOS 7 及更早版本中,如果我们想在应用程序中找到 topMostWindow,我们通常使用以下代码行 [[[UIApplication sharedApplication] windows]
我已经尝试解决这个问题很长一段时间了:我无法访问窗口的 url,因为它位于另一个域上..有一些解决方案吗? function login() { var cb = window.ope
是否可以将 FFMPEG 视频流传递到 C# 窗口?现在它在新窗口中作为新进程打开,我只是想将它传递给我自己的 SessionWindow。 此时我像这样执行ffplay: public void E
我有一个名为 x 的矩阵看起来像这样: pTime Close 1 1275087600 1.2268 2 1275264000 1.2264 3 1275264300 1.2
在编译时,发生搜索,grep搜索等,Emacs会在单独的窗口中创建一个新的缓冲区来显示结果,有没有自动跳转到那个窗口的方法?这很有用,因为我可以使用 n 和 p 而不是 M-g n 和 M-g p 移
我有一个启动 PowerShell 脚本的批处理文件。 批处理文件: START Powershell -executionpolicy RemoteSigned -noexit -file "MyS
我有一个基于菜单栏的应用程序,单击图标时会显示一个窗口。在 Mac OS X Lion 上一切正常,但由于某种原因,在 Snow Leopard 和早期版本的 Mac OS X 上会出现错误。任何时候
在 macOS 中,如何在 Xcode 和/或 Interface Builder 中创建带有“集成标题栏和工具栏”的窗口? 这是“宽标题栏”类型的窗口,已添加到 OS X 10.10 Yosemit
在浏览器 (Chrome) 中 JavaScript: var DataModler = { Data: { Something: 'value' }, Process: functi
我有 3 个 html 页面。第 1 页链接到第 2 页,第 2 页链接到第 3 页(为了简单起见)。 我希望页面 2 中的链接打开页面 3 并关闭页面 1(选项卡 1)。 据我了解,您无法使用 Ja
当点击“创建节点”按钮时,如何打开一个新的框架或窗口?我希望新框架包含一个文本字段和下拉菜单,以便用户可以选择一个选项。 Create node Search node
我有一个用户控件,用于编辑应用程序中的某些对象。 我最近遇到一个实例,我想弹出一个新的对话框(窗口)来托管此用户控件。 如何实例化新窗口并将需要设置的任何属性从窗口传递到用户控件? 感谢您的宝贵时间。
我有一个Observable,它发出许多对象,我想使用window或buffer操作对这些对象进行分组。但是,我不想指定count参数来确定窗口中应包含多少个对象,而是希望能够使用自定义条件。 例如,
我有以下代码,它打开一个新的 JavaFX 阶段(我们称之为窗口)。 openAlertBox.setOnAction(e -> { AlertBox alert = AlertBox
我要添加一个“在新窗口中打开”上下文菜单项,该菜单项将以新的UIScene打开我的应用程序文档之一。当然,我只想在实际上支持多个场景的设备上显示该菜单项。 目前,我只是在检查设备是否是使用旧设备的iP
我正在尝试创建一个 AIR 应用程序来记录应用程序的使用情况,使用 AIR 从系统获取信息的唯一简单方法是使用命令行工具和抓取 标准输出 . 我知道像 这样的工具顶部 和 ps 对于 OS X,但它们
所以我有这个简单的 turtle 螺旋制作器,我想知道是否有一种方法可以打印出由该程序创建的我的设计副本。 代码: import turtle x= float(input("Angle: ")) y
我正在编写一个 C# WPF 程序,它将文本消息发送到另一个程序的窗口。我有一个宏程序作为我的键盘驱动程序 (Logitech g15) 的一部分,它已经这样做了,尽管它不会将击键直接发送到进程,而是
我尝试使用以下代码通过 UDP 发送,但得到了奇怪的结果。 if((sendto(newSocket, sendBuf, totalLength, 0, (SOCKADDR *)&sendAd
我是一名优秀的程序员,十分优秀!