【问题标题】:Python Page Rank Streaming Application using Hadoop, py4j.protocol.Py4JJavaError: An error occurred while calling o27.partitionsPython Page Rank Streaming Application using Hadoop, py4j.protocol.Py4JJavaError: An error occurred while calling o27.partitions
【发布时间】:2020-10-25 03:25:38
【问题描述】:

我正在尝试将页面排名算法从简单的 python 代码传递到使用 python in spark 的 Streaming 应用程序,在这种情况下,它根据特定时间 (10) 秒从另一个将文件放入目录的 python 脚本中获取输入,这个脚本应该根据时间来分析它们,当我运行以下代码时出现错误,我不知道是什么原因导致错误,我试图获取 .csv 文件,

import sys
from pyspark import SparkContext
from pyspark.sql import SparkSession
from pyspark.streaming import StreamingContext

def main(input_folder_location):
sc = SparkContext.getOrCreate()
ssc = StreamingContext(sc, 10)  # Streaming will execute in each 3 seconds
ssc.checkpoint(input_folder_location)  # 'mean directory name, Directory to be checked

links = spark.sparkContext.textFile(input_folder_location). \
    map(lambda line: line.split(',')). \
    map(lambda pages: (pages[0], pages[1])). \
    distinct(). \
    groupByKey(). \
    map(lambda x: (x[0], list(x[1])))

ranks = links.map(lambda element: (element[0], 1.0))
# iterations = int(sys.argv[3])
iterations = 4

    for x in range(iterations + 1):
    contribs = links.join(ranks).flatMap(lambda row: computeContribs(row[1][0], row[1][1]))
    print("\n")
    print("------- Iter: " + str(x) + " --------")
    ranks = contribs.reduceByKey(lambda v1, v2: v1 + v2).map(lambda x: (x[0], x[1] * 0.85 + 0.15))
    for rank in ranks.collect():
        print(rank)
    print("\n")

print("------- Final Results --------")
for rank in ranks.collect():
    print(rank)

ssc.start()
ssc.awaitTermination()

def computeContribs(neighbors, rank):
for neighbor in neighbors:
    yield (neighbor, rank / len(neighbors))

if __name__ == "__main__":
if len(sys.argv) < 2:
    sys.stderr.write(
        "Error: Usage: StreamingApp.py <input-file-directory>")
    sys.exit()

spark = SparkSession.builder.getOrCreate()
spark.sparkContext.setLogLevel("WARN")

main(sys.argv[1])

【问题讨论】:

    标签: python apache-spark hadoop pyspark spark-streaming


    【解决方案1】:
    import sys
    from pyspark import SparkContext
    from pyspark.sql import SparkSession
    from pyspark.streaming import StreamingContext
    
    def computeContribs(neighbors, rank):
        for neighbor in neighbors:
            yield (neighbor, rank / len(neighbors))
    def main(input_folder_location):
    sc = SparkContext.getOrCreate()
    ssc = StreamingContext(sc, 3)   #Streaming will execute in each 3 seconds
    lines = ssc.textFileStream(input_folder_location)  #'log/ mean directory name
    print(" hello from app")
    counts = lines.map(lambda line: line.split(",")) \
    .map(lambda pages:(pages[0],pages[1])) \
    .transform(lambda rdd: rdd.distinct()) \
    .groupByKey() \
    .map(lambda x: (x[0],list(x[1])))
    
      ranks = counts.map(lambda element:(element[0],1.0))
    iterations = 5
    for x in range(1):
        contribs = counts.join(ranks).flatMap(lambda row: computeContribs(row[1][0],row[1][1]))
        print("\n")
        print(" iter --------------",x)
        ranks = contribs.reduceByKey(lambda v1,v2:v1+v2)
        print("\n")
    #counts.pprint()
    ranks.pprint()
    
    print(" finishing the task")
    ssc.start()
    ssc.awaitTermination()
    
    
    if __name__ == "__main__":
         if len(sys.argv) < 2:
            sys.stderr.write(
            "Error: Usage: StreamingApp.py <input-file-directory>")
            sys.exit()
    
    spark = SparkSession.builder.getOrCreate()
    spark.sparkContext.setLogLevel("WARN")
    
    main(sys.argv[1])
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-05-25
      • 1970-01-01
      • 2022-10-23
      • 2018-10-06
      • 1970-01-01
      • 2022-12-28
      • 2022-12-27
      • 2022-09-24
      相关资源
      最近更新 更多