【问题标题】:Python - Write to multiple outputs by key Spark - one Spark jobPython - 通过键写入多个输出 Spark - 一个 Spark 作业
【发布时间】:2015-10-19 08:14:20
【问题描述】:

如何在一项作业中使用 Python 和 Spark 为 RDD 中的每个键写入多个输出?我知道我可以尝试对所有可能的键使用 .filter,但这是很多工作,会创造很多工作。

类似于这个问题: Write to multiple outputs by key Spark - one Spark job

但是,上述问题的答案是在 scala 中。寻找如何使用 Python。

PATH = os.path.join("s3://asdf/hjkl", 'temp_date', "intermediate_data/")
global current_sport
current_sport = ''
def format_for_output(x):
    current_sport = x[0]
    return json.dumps(x[1])
recommendation2.map(format_for_output).saveAsTextFile(os.path.join(PATH, current_sport))

【问题讨论】:

  • API 很相似,应该很容易转译。

标签: python apache-spark pyspark


【解决方案1】:

如果您想要简单的 Python 解决方案,那么您可以简单地按键分区 RDD。首先让我们创建一些虚拟数据:

import numpy as np
np.random.seed(1)

keys = [chr(x) for x in xrange(65, 91)]
rdd = sc.parallelize(
    (np.random.choice(keys), np.random.randint(0, 100)) for _ in xrange(10000))

现在让我们假设我们对密钥一无所知。我们必须创建从键到分区 id 的映射:

mapping = sc.broadcast(
    rdd.keys(). # Get keys
        distinct(). # Find unique
        sortBy(lambda x: x). # Sort
        zipWithIndex(). # Add index
        collectAsMap()) # Create dict

最后我们可以使用上面的映射进行分区并保存到文本文件:

(rdd.
    partitionBy(
        len(mapping.value) # Number of partitions
        partitionFunc=lambda x: mapping.value.get(x) # Mapping
    ).saveAsTextFile("foo"))

让我们检查一切是否按预期工作:

import glob

cnts = rdd.countByKey() # Count values by key
fs = sorted(glob.glob("foo/part-*")) # Get output names

assert len(fs) == len(mapping.value) # All keys present

for (k, v) in sorted(mapping.value.items()):
    with open(fs[v]) as fr:
        lines = fr.readlines()
        assert len(lines) == cnts[k] # Number of records as expected
        assert all(k in line for line in lines) # All with the same key

【讨论】:

    猜你喜欢
    • 2014-07-22
    • 1970-01-01
    • 1970-01-01
    • 2021-09-18
    • 2018-12-12
    • 1970-01-01
    相关资源
    最近更新 更多