【发布时间】:2016-06-15 12:08:59
【问题描述】:
根据标题。我知道textFile,但顾名思义,它只适用于文本文件。
我需要访问 HDFS 或本地路径上的路径内的文件/目录。我正在使用 pyspark。
【问题讨论】:
标签: hadoop apache-spark pyspark
根据标题。我知道textFile,但顾名思义,它只适用于文本文件。
我需要访问 HDFS 或本地路径上的路径内的文件/目录。我正在使用 pyspark。
【问题讨论】:
标签: hadoop apache-spark pyspark
使用 JVM 网关可能不是那么优雅,但在某些情况下,下面的代码可能会有所帮助:
URI = sc._gateway.jvm.java.net.URI
Path = sc._gateway.jvm.org.apache.hadoop.fs.Path
FileSystem = sc._gateway.jvm.org.apache.hadoop.fs.FileSystem
Configuration = sc._gateway.jvm.org.apache.hadoop.conf.Configuration
fs = FileSystem.get(URI("hdfs://somehost:8020"), Configuration())
status = fs.listStatus(Path('/some_dir/yet_another_one_dir/'))
for fileStatus in status:
print(fileStatus.getPath())
【讨论】:
globStatus 而不是 fileStatus,例如status = fs.globStatus(Path('/some_dir/yet_another_one_dir/*.csv'))
somehost,即namenode的好方法是什么?
files = [file.getPath() for file in status] 需要一段时间。正常吗?我会说这不是最有效的方法。
globStatus 而不是listStatus,而不是fileStatus(这只是一个临时变量)。
我认为将 Spark 仅视为一种数据处理工具会很有帮助,它的域从加载数据开始。它可以读取多种格式,并且支持 Hadoop glob 表达式,这对于从 HDFS 中的多个路径读取非常有用,但它没有我知道的用于遍历目录或文件的内置工具,也没有专用于与 Hadoop 或 HDFS 交互的实用程序。
有一些可用的工具可以满足您的需求,包括esutil 和hdfs。 hdfs lib 支持 CLI 和 API,您可以直接跳转到“我如何在 Python 中列出 HDFS 文件”here。它看起来像这样:
from hdfs import Config
client = Config().get_client('dev')
files = client.list('the_dir_path')
【讨论】:
default.alias
如果你使用PySpark,你可以交互式地执行命令:
列出所选目录中的所有文件:
hdfs dfs -ls <path> 例如:hdfs dfs -ls /user/path:
import os
import subprocess
cmd = 'hdfs dfs -ls /user/path'
files = subprocess.check_output(cmd, shell=True).strip().split('\n')
for path in files:
print path
或在所选目录中搜索文件:
hdfs dfs -find <path> -name <expression> 例如:hdfs dfs -find /user/path -name *.txt:
import os
import subprocess
cmd = 'hdfs dfs -find {} -name *.txt'.format(source_dir)
files = subprocess.check_output(cmd, shell=True).strip().split('\n')
for path in files:
filename = path.split(os.path.sep)[-1].split('.txt')[0]
print path, filename
【讨论】:
hdfs dfs -rm -r 命令?是使用相同的 check_output 方法还是其他方式?
这可能对你有用:
import subprocess, re
def listdir(path):
files = str(subprocess.check_output('hdfs dfs -ls ' + path, shell=True))
return [re.search(' (/.+)', i).group(1) for i in str(files).split("\\n") if re.search(' (/.+)', i)]
listdir('/user/')
这也有效:
hadoop = sc._jvm.org.apache.hadoop
fs = hadoop.fs.FileSystem
conf = hadoop.conf.Configuration()
path = hadoop.fs.Path('/user/')
[str(f.getPath()) for f in fs.get(conf).listStatus(path)]
【讨论】:
使用snakebite 库有一个简单的方法来做到这一点
from snakebite.client import Client
hadoop_client = Client(HADOOP_HOST, HADOOP_PORT, use_trash=False)
for x in hadoop_client.ls(['/']):
... print x
【讨论】: