【问题标题】:Python read file as stream from HDFSPython从HDFS读取文件作为流
【发布时间】:2012-09-11 05:32:50
【问题描述】:

这是我的问题:我在 HDFS 中有一个可能很大的文件(=不足以容纳所有内存)

我想做的是避免将此文件缓存在内存中,而只像处理常规文件一样逐行处理它:

for line in open("myfile", "r"):
    # do some processing

我正在寻找是否有一种简单的方法可以在不使用外部库的情况下正确完成这项工作。我可能可以使它与libpyhdfspython-hdfs 一起工作,但如果可能的话,我希望避免在系统中引入新的依赖项和未经测试的库,特别是因为这两个似乎都没有得到大量维护并声明它们应该不能用于生产。

我正在考虑使用标准的“hadoop”命令行工具使用 Python subprocess 模块来执行此操作,但我似乎无法执行我需要的操作,因为没有命令行工具可以执行此操作我的处理,我想以流方式为每一行执行一个 Python 函数。

有没有办法使用 subprocess 模块将 Python 函数应用为管道的正确操作数?或者更好的是,像文件一样打开它作为生成器,这样我就可以轻松处理每一行?

cat = subprocess.Popen(["hadoop", "fs", "-cat", "/path/to/myfile"], stdout=subprocess.PIPE)

如果有另一种方法可以在不使用外部库的情况下实现我上面描述的内容,我也很开放。

感谢您的帮助!

【问题讨论】:

标签: python hadoop subprocess hdfs


【解决方案1】:

你想要xreadlines,它从文件中读取行而不将整个文件加载到内存中。

编辑

现在我看到了你的问题,你只需要从 Popen 对象中获取标准输出管道:

cat = subprocess.Popen(["hadoop", "fs", "-cat", "/path/to/myfile"], stdout=subprocess.PIPE)
for line in cat.stdout:
    print line

【讨论】:

  • 这与for line in open("myfile") 有何不同?我认为这里的困难在于我不是在处理常规文件,而是在 HDFS 中的文件,我想知道是否有一种方法可以让子进程模块(或其他东西)逐行处理它一些 Python 代码。 (我没有投反对票)
  • 我想我当时误解了你的问题。你不想在python中逐行处理它吗? HDFS 没有办法cat 一个文件吗? (我希望如此。)只需将其称为子进程并将结果包装在 xreadlinesfor line in... 中。
  • 请注意,如果将 -cat 替换为 -text,它也会处理压缩。
  • 请注意,从 2.3 开始,xreadlines is deprecated(只需使用for line in file,就像在您的编辑中一样)。
  • @CharlesMenguy:在 Python 2 上,您可以使用 for line in iter(cat.stdout.readline, ''): print line,(注意:末尾的逗号):它解决了两个问题的答案:1. Python 2 中存在预读错误这会延迟 cat 命令的输出 2. print line(无逗号)可能会引入不必要的换行符
【解决方案2】:

如果您想不惜一切代价避免添加外部依赖项,Keith 的答案就是要走的路。另一方面,Pydoop 可以让您的生活更轻松:

import pydoop.hdfs as hdfs
with hdfs.open('/user/myuser/filename') as f:
    for line in f:
        do_something(line)

关于您的顾虑,Pydoop 正在积极开发中,并已在 CRS4 用于生产多年,主要用于计算生物学应用。

西蒙娜

【讨论】:

  • hdfscli 有类似的功能,而且更轻量级。不要忘记启用 WebHDFS 以使用它。
  • @simleo 我得到一个 AttributeError: 'module' object has no attribute open
  • @JohnConstantine 如果您在使用 Pydoop 时遇到问题,您应该在 github.com/crs4/pydoop/issues 提出问题
【解决方案3】:

在过去的两年里,Hadoop-Streaming 发生了很多变化。根据 Cloudera 的说法,这非常快:http://blog.cloudera.com/blog/2013/01/a-guide-to-python-frameworks-for-hadoop/ 我已经取得了很好的成功。

【讨论】:

    【解决方案4】:

    您可以使用 WebHDFS Python 库(基于 urllib3 构建):

    from hdfs import InsecureClient
    client_hdfs = InsecureClient('http://host:port', user='root')
    with client_hdfs.write(access_path) as writer:
        dump(records, writer)  # tested for pickle and json (doesnt work for joblib)
    

    或者你可以在python中使用requests包:

    import requests
    from json import dumps
    params = (('op', 'CREATE')
    ('buffersize', 256))
    data = dumps(file)  # some file or object - also tested for pickle library
    response = requests.put('http://host:port/path', params=params, data=data)  # response 200 = successful
    

    希望这会有所帮助!

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-06-14
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-09-04
      • 1970-01-01
      • 2021-06-18
      • 2015-05-11
      相关资源
      最近更新 更多