【发布时间】: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