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

windowns使用PySpark環(huán)境配置和基本操作

 更新時(shí)間:2021年05月17日 12:05:04   作者:Nick_Spider  
pyspark是Spark對(duì)Python的api接口,可以在Python環(huán)境中通過調(diào)用pyspark模塊來操作spark,這篇文章主要介紹了windowns使用PySpark環(huán)境配置和基本操作,感興趣的可以了解一下

下載依賴

首先需要下載hadoop和spark,解壓,然后設(shè)置環(huán)境變量。
hadoop清華源下載
spark清華源下載

HADOOP_HOME => /path/hadoop
SPARK_HOME => /path/spark

安裝pyspark。

pip install pyspark

基本使用

可以在shell終端,輸入pyspark,有如下回顯:

輸入以下指令進(jìn)行測(cè)試,并創(chuàng)建SparkContext,SparkContext是任何spark功能的入口點(diǎn)。

>>> from pyspark import SparkContext
>>> sc = SparkContext("local", "First App")

如果以上不會(huì)報(bào)錯(cuò),恭喜可以開始使用pyspark編寫代碼了。
不過,我這里使用IDE來編寫代碼,首先我們先在終端執(zhí)行以下代碼關(guān)閉SparkContext。

>>> sc.stop()

下面使用pycharm編寫代碼,如果修改了環(huán)境變量需要先重啟pycharm。
在pycharm運(yùn)行如下程序,程序會(huì)起本地模式的spark計(jì)算引擎,通過spark統(tǒng)計(jì)abc.txt文件中a和b出現(xiàn)行的數(shù)量,文件路徑需要自己指定。

from pyspark import SparkContext

sc = SparkContext("local", "First App")
logFile = "abc.txt"
logData = sc.textFile(logFile).cache()
numAs = logData.filter(lambda s: 'a' in s).count()
numBs = logData.filter(lambda s: 'b' in s).count()
print("Line with a:%i,line with b:%i" % (numAs, numBs))

運(yùn)行結(jié)果如下:

20/03/11 16:15:57 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
20/03/11 16:15:58 WARN Utils: Service 'SparkUI' could not bind on port 4040. Attempting port 4041.
Line with a:3,line with b:1

這里說一下,同樣的工作使用python可以做,spark也可以做,使用spark主要是為了高效的進(jìn)行分布式計(jì)算。
戳pyspark教程
戳spark教程

RDD

RDD代表Resilient Distributed Dataset,它們是在多個(gè)節(jié)點(diǎn)上運(yùn)行和操作以在集群上進(jìn)行并行處理的元素,RDD是spark計(jì)算的操作對(duì)象。
一般,我們先使用數(shù)據(jù)創(chuàng)建RDD,然后對(duì)RDD進(jìn)行操作。
對(duì)RDD操作有兩種方法:
Transformation(轉(zhuǎn)換) - 這些操作應(yīng)用于RDD以創(chuàng)建新的RDD。例如filter,groupBy和map。
Action(操作) - 這些是應(yīng)用于RDD的操作,它指示Spark執(zhí)行計(jì)算并將結(jié)果發(fā)送回驅(qū)動(dòng)程序,例如count,collect等。

創(chuàng)建RDD

parallelize是從列表創(chuàng)建RDD,先看一個(gè)例子:

from pyspark import SparkContext


sc = SparkContext("local", "count app")
words = sc.parallelize(
    ["scala",
     "java",
     "hadoop",
     "spark",
     "akka",
     "spark vs hadoop",
     "pyspark",
     "pyspark and spark"
     ])
print(words)

結(jié)果中我們得到一個(gè)對(duì)象,就是我們列表數(shù)據(jù)的RDD對(duì)象,spark之后可以對(duì)他進(jìn)行操作。

ParallelCollectionRDD[0] at parallelize at PythonRDD.scala:195

Count

count方法返回RDD中的元素個(gè)數(shù)。

from pyspark import SparkContext


sc = SparkContext("local", "count app")
words = sc.parallelize(
    ["scala",
     "java",
     "hadoop",
     "spark",
     "akka",
     "spark vs hadoop",
     "pyspark",
     "pyspark and spark"
     ])
print(words)

counts = words.count()
print("Number of elements in RDD -> %i" % counts)

返回結(jié)果:

Number of elements in RDD -> 8

Collect

collect返回RDD中的所有元素。

from pyspark import SparkContext


sc = SparkContext("local", "collect app")
words = sc.parallelize(
    ["scala",
     "java",
     "hadoop",
     "spark",
     "akka",
     "spark vs hadoop",
     "pyspark",
     "pyspark and spark"
     ])
coll = words.collect()
print("Elements in RDD -> %s" % coll)

返回結(jié)果:

Elements in RDD -> ['scala', 'java', 'hadoop', 'spark', 'akka', 'spark vs hadoop', 'pyspark', 'pyspark and spark']

foreach

每個(gè)元素會(huì)使用foreach內(nèi)的函數(shù)進(jìn)行處理,但是不會(huì)返回任何對(duì)象。
下面的程序中,我們定義的一個(gè)累加器accumulator,用于儲(chǔ)存在foreach執(zhí)行過程中的值。

from pyspark import SparkContext
sc = SparkContext("local", "ForEach app")

accum = sc.accumulator(0)
data = [1, 2, 3, 4, 5]
rdd = sc.parallelize(data)


def increment_counter(x):
    print(x)
    accum.add(x)
 return 0

s = rdd.foreach(increment_counter)
print(s)  # None
print("Counter value: ", accum)

返回結(jié)果:

None
Counter value:  15

filter

返回一個(gè)包含元素的新RDD,滿足過濾器的條件。

from pyspark import SparkContext
sc = SparkContext("local", "Filter app")
words = sc.parallelize(
    ["scala",
     "java",
     "hadoop",
     "spark",
     "akka",
     "spark vs hadoop",
     "pyspark",
     "pyspark and spark"]
)
words_filter = words.filter(lambda x: 'spark' in x)
filtered = words_filter.collect()
print("Fitered RDD -> %s" % (filtered))

 

Fitered RDD -> ['spark', 'spark vs hadoop', 'pyspark', 'pyspark and spark']

也可以改寫成這樣:

from pyspark import SparkContext
sc = SparkContext("local", "Filter app")
words = sc.parallelize(
    ["scala",
     "java",
     "hadoop",
     "spark",
     "akka",
     "spark vs hadoop",
     "pyspark",
     "pyspark and spark"]
)


def g(x):
    for i in x:
        if "spark" in x:
            return i

words_filter = words.filter(g)
filtered = words_filter.collect()
print("Fitered RDD -> %s" % (filtered))

map

將函數(shù)應(yīng)用于RDD中的每個(gè)元素并返回新的RDD。

from pyspark import SparkContext
sc = SparkContext("local", "Map app")
words = sc.parallelize(
    ["scala",
     "java",
     "hadoop",
     "spark",
     "akka",
     "spark vs hadoop",
     "pyspark",
     "pyspark and spark"]
)
words_map = words.map(lambda x: (x, 1, "_{}".format(x)))
mapping = words_map.collect()
print("Key value pair -> %s" % (mapping))

返回結(jié)果:

Key value pair -> [('scala', 1, '_scala'), ('java', 1, '_java'), ('hadoop', 1, '_hadoop'), ('spark', 1, '_spark'), ('akka', 1, '_akka'), ('spark vs hadoop', 1, '_spark vs hadoop'), ('pyspark', 1, '_pyspark'), ('pyspark and spark', 1, '_pyspark and spark')]

Reduce

執(zhí)行指定的可交換和關(guān)聯(lián)二元操作后,然后返回RDD中的元素。

from pyspark import SparkContext
from operator import add


sc = SparkContext("local", "Reduce app")
nums = sc.parallelize([1, 2, 3, 4, 5])
adding = nums.reduce(add)
print("Adding all the elements -> %i" % (adding))

 這里的add是python內(nèi)置的函數(shù),可以使用ide查看:

def add(a, b):
    "Same as a + b."
    return a + b

reduce會(huì)依次對(duì)元素相加,相加后的結(jié)果加上其他元素,最后返回結(jié)果(RDD中的元素)。

Adding all the elements -> 15

Join

返回RDD,包含兩者同時(shí)匹配的鍵,鍵包含對(duì)應(yīng)的所有元素。

from pyspark import SparkContext


sc = SparkContext("local", "Join app")
x = sc.parallelize([("spark", 1), ("hadoop", 4), ("python", 4)])
y = sc.parallelize([("spark", 2), ("hadoop", 5)])
print("x =>", x.collect())
print("y =>", y.collect())
joined = x.join(y)
final = joined.collect()
print( "Join RDD -> %s" % (final))

返回結(jié)果:

x => [('spark', 1), ('hadoop', 4), ('python', 4)]
y => [('spark', 2), ('hadoop', 5)]
Join RDD -> [('hadoop', (4, 5)), ('spark', (1, 2))]

到此這篇關(guān)于windowns使用PySpark環(huán)境配置和基本操作的文章就介紹到這了,更多相關(guān)PySpark環(huán)境配置 內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Django之第三方平臺(tái)QQ授權(quán)登錄的實(shí)現(xiàn)

    Django之第三方平臺(tái)QQ授權(quán)登錄的實(shí)現(xiàn)

    本文主要介紹了Django之第三方平臺(tái)QQ授權(quán)登錄的實(shí)現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2023-05-05
  • 五分鐘學(xué)會(huì)Python 模塊和包、文件

    五分鐘學(xué)會(huì)Python 模塊和包、文件

    通過學(xué)習(xí)本文可以五分鐘掌握Python 模塊和包、文件的相關(guān)知識(shí),本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-08-08
  • 搭建pypi私有倉庫實(shí)現(xiàn)過程詳解

    搭建pypi私有倉庫實(shí)現(xiàn)過程詳解

    這篇文章主要介紹了搭建pypi私有倉庫實(shí)現(xiàn)過程詳解,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-11-11
  • scipy.interpolate插值方法實(shí)例講解

    scipy.interpolate插值方法實(shí)例講解

    這篇文章主要介紹了scipy.interpolate插值方法介紹,本文結(jié)合實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-12-12
  • 利用Python對(duì)中國500強(qiáng)排行榜數(shù)據(jù)進(jìn)行可視化分析

    利用Python對(duì)中國500強(qiáng)排行榜數(shù)據(jù)進(jìn)行可視化分析

    這篇文章主要介紹了利用Python對(duì)中國500強(qiáng)排行榜數(shù)據(jù)進(jìn)行可視化分析,從不同角度去對(duì)數(shù)據(jù)進(jìn)行統(tǒng)計(jì)分析,可視化展示,下文詳細(xì)內(nèi)容介紹需要的小伙伴可以參考一下
    2022-05-05
  • Python進(jìn)行密碼學(xué)反向密碼教程

    Python進(jìn)行密碼學(xué)反向密碼教程

    這篇文章主要為大家介紹了Python進(jìn)行密碼學(xué)反向密碼的教程詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-05-05
  • Django Rest framework權(quán)限的詳細(xì)用法

    Django Rest framework權(quán)限的詳細(xì)用法

    這篇文章主要介紹了Django Rest framework權(quán)限的詳細(xì)用法,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-07-07
  • python中獲得當(dāng)前目錄和上級(jí)目錄的實(shí)現(xiàn)方法

    python中獲得當(dāng)前目錄和上級(jí)目錄的實(shí)現(xiàn)方法

    這篇文章主要介紹了python中獲得當(dāng)前目錄和上級(jí)目錄的實(shí)現(xiàn)方法,需要的朋友可以參考下
    2017-10-10
  • python assert斷言的實(shí)例用法

    python assert斷言的實(shí)例用法

    在本篇文章里小編給大家整理了一篇關(guān)于python assert斷言的實(shí)例用法,有需要的朋友們可以跟著學(xué)習(xí)參考下。
    2021-09-09
  • python  logging日志打印過程解析

    python logging日志打印過程解析

    這篇文章主要介紹了python logging日志打印過程解析,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-10-10

最新評(píng)論

宜丰县| 贵阳市| 翼城县| 娱乐| 唐山市| 集安市| 阳泉市| 岢岚县| 抚松县| 大兴区| 东港市| 淳安县| 依兰县| 麦盖提县| 巴楚县| 忻城县| 沁水县| 喜德县| 荆州市| 缙云县| 金平| 麻栗坡县| 德钦县| 亳州市| 仁怀市| 巴彦淖尔市| 林周县| 方正县| 江永县| 抚松县| 大兴区| 武义县| 民乐县| 舟山市| 武冈市| 延庆县| 新沂市| 惠水县| 兴和县| 南康市| 长岛县|