最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

spark?dataframe全局排序id與分組后保留最大值行

 更新時(shí)間:2023年02月09日 08:49:25   作者:算法全棧之路  
這篇文章主要為大家介紹了spark?dataframe全局排序id與分組后保留最大值行實(shí)現(xiàn)詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

正文

作為一個(gè)算法工程師,日常學(xué)習(xí)和工作中,不光要 訓(xùn)練模型關(guān)注效果 ,更多的 時(shí)間 是在 準(zhǔn)備樣本數(shù)據(jù)與分析數(shù)據(jù) 等,而這些過(guò)程 都與 大數(shù)據(jù) spark和hadoop生態(tài) 的若干工具息息相關(guān)。

今天我們就不在更新 機(jī)器學(xué)習(xí)算法模型 相關(guān)的內(nèi)容,分享兩個(gè) spark函數(shù) 吧,以前也在某種場(chǎng)景中使用過(guò)但沒(méi)有保存收藏,哎??! 事前不搜藏,臨時(shí)抱佛腳 的感覺(jué) 真是 痛苦,太耽誤干活了 。

so,把這 兩個(gè)函數(shù) 記在這里 以備不時(shí) 之需~

(1) 得到 spark dataframe 全局排序ID

這個(gè)函數(shù)的 應(yīng)用場(chǎng)景 就是:根據(jù)某一列的數(shù)值對(duì) spark 的 dataframe 進(jìn)行排序, 得到全局多分區(qū)排序的全局有序ID,新增一列保存這個(gè)rank id ,并且保留別的列的數(shù)據(jù)無(wú)變化 。

有用戶會(huì)說(shuō),這不是很容易嗎 ,直接用 orderBy 不就可以了嗎,但是難點(diǎn)是:orderBy完記錄下全局ID 并且 保持原來(lái)全部列的DF數(shù)據(jù) 。

多說(shuō)無(wú)益,遇到這個(gè)場(chǎng)景 直接copy 用起來(lái) 就知道 有多爽 了,同類(lèi)問(wèn)題 我們可以 用下面 這個(gè)函數(shù) 解決 ~

scala 寫(xiě)的 spark 版本代碼:

def dfZipWithIndex(
  df: DataFrame,
  offset: Int = 1,
  colName: String ="rank_id",
  inFront: Boolean = true
) : DataFrame = {
  df.sqlContext.createDataFrame(
    df.rdd.zipWithIndex.map(ln =>
      Row.fromSeq(
        (if (inFront) Seq(ln._2 + offset) else Seq())
          ++ ln._1.toSeq ++
        (if (inFront) Seq() else Seq(ln._2 + offset))
      )
    ),
    StructType(
      (if (inFront) Array(StructField(colName,LongType,false)) else Array[StructField]())
        ++ df.schema.fields ++
      (if (inFront) Array[StructField]() else Array(StructField(colName,LongType,false)))
    )
  )
}

函數(shù)調(diào)用我們可以用這行代碼調(diào)用: val ranked_df = dfZipWithIndex(raw_df.orderBy($"predict_score".desc)), 直接復(fù)制過(guò)去就可以~

python寫(xiě)的 pyspark 版本代碼:

from pyspark.sql.types import LongType, StructField, StructType
def dfZipWithIndex (df, offset=1, colName="rank_id"):
    new_schema = StructType(
                    [StructField(colName,LongType(),True)]        # new added field in front
                    + df.schema.fields                            # previous schema
                )
    zipped_rdd = df.rdd.zipWithIndex()
    new_rdd = zipped_rdd.map(lambda (row,rowId): ([rowId +offset] + list(row)))
    return spark.createDataFrame(new_rdd, new_schema)

調(diào)用 同理 , 這里我就不在進(jìn)行贅述了。

(2)分組后保留最大值行

這個(gè)函數(shù)的 應(yīng)用場(chǎng)景 就是: 當(dāng)我們使用 spark 或則 sparkSQL 查找某個(gè) dataframe 數(shù)據(jù)的時(shí)候,在某一天里,任意一個(gè)用戶可能有多條記錄,我們需要 對(duì)每一個(gè)用戶,保留dataframe 中 某列值最大 的那行數(shù)據(jù) 。

其中的 關(guān)鍵點(diǎn) 在于:一次性求出對(duì)每個(gè)用戶分組后,求得每個(gè)用戶的多行記錄中,某個(gè)值最大的行進(jìn)行數(shù)據(jù)保留 。

當(dāng)然,經(jīng)過(guò) 簡(jiǎn)單修改代碼,不一定是最大,最小也是可以的,平均都o(jì)k

scala 寫(xiě)的 spark 版本代碼:

// 得到一天內(nèi)一個(gè)用戶多個(gè)記錄里面時(shí)間最大的那行用戶的記錄
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions
val w = Window.partitionBy("user_id")
val result_df = raw_df
    .withColumn("max_time",functions.max("time").over(w))
    .where($"time" === $"max_time")
    .drop($"max_time")

python寫(xiě)的 pyspark 版本代碼:

# pyspark dataframe 某列值最大的元素所在的那一行 
# GroupBy 列并過(guò)濾 Pyspark 中某列值最大的行 
# 創(chuàng)建一個(gè)Window 以按A列進(jìn)行分區(qū),并使用它來(lái)計(jì)算每個(gè)組的最大值。然后過(guò)濾出行,使 B 列中的值等于最大值 
from pyspark.sql import Window
w = Window.partitionBy('user_id')
result_df = spark.sql(raw_df).withColumn('max_time', fun.max('time').over(w))\
    .where(fun.col('time') == fun.col('time'))
    .drop('max_time')

我們可以看到: 這個(gè)函數(shù)的關(guān)鍵就是運(yùn)用了 spark 的 window 函數(shù) ,靈活運(yùn)用 威力無(wú)窮 哦 !

到這里,spark利器2函數(shù)之dataframe全局排序id與分組后保留最大值行 的全文 就寫(xiě)完了 ,更多關(guān)于spark dataframe全局排序的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • 淺談keras保存模型中的save()和save_weights()區(qū)別

    淺談keras保存模型中的save()和save_weights()區(qū)別

    這篇文章主要介紹了淺談keras保存模型中的save()和save_weights()區(qū)別,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-05-05
  • Python實(shí)現(xiàn)將Excel轉(zhuǎn)換為json的方法示例

    Python實(shí)現(xiàn)將Excel轉(zhuǎn)換為json的方法示例

    這篇文章主要介紹了Python實(shí)現(xiàn)將Excel轉(zhuǎn)換為json的方法,涉及Python文件讀寫(xiě)及格式轉(zhuǎn)換相關(guān)操作技巧,需要的朋友可以參考下
    2017-08-08
  • Python中的localtime()方法使用詳解

    Python中的localtime()方法使用詳解

    這篇文章主要介紹了Python中的localtime()方法使用詳解,是Python入門(mén)學(xué)習(xí)的基礎(chǔ)知識(shí),需要的朋友可以參考下
    2015-05-05
  • python實(shí)現(xiàn)微信小程序的多種支付方式

    python實(shí)現(xiàn)微信小程序的多種支付方式

    這篇文章主要為大家介紹了python實(shí)現(xiàn)微信小程序的多種支付方式的實(shí)現(xiàn)示例,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步早日升職加薪
    2022-04-04
  • Python3+Pygame實(shí)現(xiàn)射擊游戲完整代碼

    Python3+Pygame實(shí)現(xiàn)射擊游戲完整代碼

    這篇文章主要介紹了Python3+Pygame實(shí)現(xiàn)射擊游戲完整代碼,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-03-03
  • 詳解Python pygame安裝過(guò)程筆記

    詳解Python pygame安裝過(guò)程筆記

    本篇文章主要介紹了詳解Python pygame安裝過(guò)程筆記。小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2017-06-06
  • Python實(shí)現(xiàn)壓縮pdf文件大小

    Python實(shí)現(xiàn)壓縮pdf文件大小

    工作中常需要壓縮數(shù)據(jù)文件大小,壓縮PDF文件是一種減少PDF文件大小的方法,這樣可以使文件更易于傳輸和存儲(chǔ),本文將使用Python實(shí)現(xiàn)這一功能,需要的可以參考下
    2024-02-02
  • 如何獲取numpy array前N個(gè)最大值

    如何獲取numpy array前N個(gè)最大值

    這篇文章主要介紹了獲取numpy array前N個(gè)最大值的操作,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-05-05
  • Python?NumPy教程之?dāng)?shù)組的基本操作詳解

    Python?NumPy教程之?dāng)?shù)組的基本操作詳解

    Numpy?中的數(shù)組是一個(gè)元素表(通常是數(shù)字),所有元素類(lèi)型相同,由正整數(shù)元組索引。本文將通過(guò)一些示例詳細(xì)講一下NumPy中數(shù)組的一些基本操作,需要的可以參考一下
    2022-08-08
  • Python生成requirements.txt的三種方法

    Python生成requirements.txt的三種方法

    requirements.txt?文件通常用于列出項(xiàng)目所需的所有Python包及其版本,本文主要介紹了Python生成requirements.txt的三種方法,具有一定的參考價(jià)值,感興趣的可以了解一下
    2024-07-07

最新評(píng)論

左权县| 巴楚县| 广宗县| 鹿邑县| 遂宁市| 泉州市| 天门市| 金寨县| 丰宁| 阿坝| 揭西县| 高阳县| 岳阳市| 兴和县| 密云县| 沈丘县| 延庆县| 揭东县| 汶川县| 嘉义市| 闻喜县| 樟树市| 英吉沙县| 泸州市| 滦南县| 台安县| 钟祥市| 甘洛县| 治多县| 尼勒克县| 长阳| 新和县| 青河县| 休宁县| 绥芬河市| 弋阳县| 富锦市| 太原市| 精河县| 措美县| 万全县|