【问题标题】:Spark - missing 1 required position argument (lambda function)Spark - 缺少 1 个必需的位置参数(lambda 函数)
【发布时间】:2018-01-08 10:55:21
【问题描述】:

我正在尝试使用 Spark 在多个服务器之间分发一些从 PDF 中提取的文本。这是使用我制作的自定义 Python 模块,是implementation of this question。 “extractTextFromPdf”函数有 2 个参数:一个表示文件路径的字符串,以及一个用于确定各种提取约束的配置文件。在这种情况下,配置文件只是一个简单的 YAML 文件,与运行提取的 Python 脚本位于同一文件夹中,并且这些文件只是在 Spark 服务器之间复制。

我遇到的主要问题是能够使用文件名而不是文件的内容作为第一个参数来调用我的提取函数。这是我目前拥有的基本脚本,在 files 文件夹中的 2 个 PDF 上运行它:

#!/usr/bin/env python3

import ScannedTextExtractor.STE as STE

from pyspark import SparkContext
sc = SparkContext("local", "STE")

input = sc.binaryFiles("/home/ubuntu/files")
processed = input.map(lambda filename, content: (STE.extractTextFromPdf(filename,'ste-config.yaml'), content))

print("Results:")
print(processed.take(2))

这会产生 lambda 错误 Missing 1 position argument: 'content'。我并不真正关心使用 PDF 的原始内容,因为我的提取函数的参数只是 PDF 的路径,而不是实际的 PDF 内容本身,我试图只给 lambda 函数提供 1 个参数。例如

processed = input.map(lambda filename: STE.extractTextFromPdf(filename,'ste-config.yaml'))

但是我遇到了问题,因为使用此设置 Spark 将 PDF 内容(作为字节流)设置为这个单数参数,但我的模块需要一个带有 PDF 路径的字符串作为第一个参数,而不是整个字节内容PDF。

我打印了 SparkContext 加载的二进制文件的 RDD,我可以看到 RDD 中有文件名和文件内容(PDF 的字节流)。但是如何将它与需要以下语法的自定义 Python 模块一起使用:

STE.extractTextFromPDF('/path/to/pdf','/path/to/config-file')

我已经尝试了 lambda 函数的多种排列,我已经三次检查了 Spark 的 RDD 和 SparkContext API。我似乎无法让它工作。

【问题讨论】:

    标签: python apache-spark lambda pyspark rdd


    【解决方案1】:

    如果你只想要路径而不是内容,那么你不应该使用sc.binaryFiles。在这种情况下,您应该并行化路径,然后让 Python 代码单独加载每个文件,如下所示:

    paths = ['/path/to/file1', '/path/to/file2']
    input = sc.parallelize(paths)
    processed = input.map(lambda path: (path, processFile(path)))
    

    这当然假设每个执行程序 Python 进程都可以直接访问文件。例如,这不适用于 HDFS 或 S3。你的库不能直接获取二进制内容吗?

    【讨论】:

    • 好吧,我试试看。不幸的是,不,它最初不是这样设计的,只是使用路径。最终我会添加那个选项
    【解决方案2】:

    map 将函数作为单个参数并传递两个参数的函数:

     input.map(lambda filename, content: (STE.extractTextFromPdf(filename,'ste-config.yaml'), content)
    

    使用任一

    input.map(lambda fc: (STE.extractTextFromPdf(fc[0],'ste-config.yaml'), fc[1])
    

    def process(x):
        filename, content = x
        return STE.extractTextFromPdf(filename,'ste-config.yaml'), content
    

    并不是说它通常不会起作用,除非:

    • STE.extractTextFromPdf 可以使用 Hadoop 兼容的文件系统或
    • 您使用 POSIX 兼容的文件系统。

    如果不行你可以试试:

    • 使用io.BytesIO 之类的伪文件(如果它支持在某个级别从类似文件的对象中读取)。
    • content 写入本地FS 上的临时文件并从那里读取。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2014-09-13
      • 1970-01-01
      • 2023-04-01
      • 1970-01-01
      • 2019-10-10
      • 2017-03-27
      相关资源
      最近更新 更多