这是一个端到端的示例,它调整了上一个问题中提到的技术,以更贴近您的问题。
Python read file as stream from HDFS
这是一个小型 Python Hadoop 流应用程序,它读取键值对,根据存储在 HDFS 中的 XML 配置文件检查键,然后仅当键与配置匹配时才发出值。匹配逻辑被卸载到一个单独的 Process.py 模块中,该模块通过使用对 hdfs dfs -cat 的外部调用从 HDFS 读取 XML 配置文件。
首先,我们创建一个名为 pythonapp 的目录,其中包含用于我们实现的 Python 源文件。稍后我们会在提交流式作业时看到,我们将在 -files 参数中传递此目录。
为什么我们将文件放入中间目录而不是在-files 参数中单独列出每个文件?这是因为当 YARN 本地化文件以在容器中执行时,它引入了一个符号链接间接层。然后 Python 无法通过符号链接正确加载模块。解决方案是将两个文件打包到同一个目录中。然后,当 YARN 本地化文件时,符号链接间接在目录级别完成,而不是在单个文件。由于主脚本和模块在物理上位于同一目录中,因此 Python 将能够正确加载模块。这个问题更详细地解释了这个问题:
How to import a custom module in a MapReduce job?
映射器.py
import subprocess
import sys
from Process import match
for line in sys.stdin:
key, value = line.split()
if match(key):
print value
进程.py
import subprocess
import xml.etree.ElementTree as ElementTree
hdfsCatProcess = subprocess.Popen(
['hdfs', 'dfs', '-cat', '/pythonAppConf.xml'],
stdout=subprocess.PIPE)
pythonAppConfXmlTree = ElementTree.parse(hdfsCatProcess.stdout)
matchString = pythonAppConfXmlTree.find('./matchString').text.strip()
def match(key):
return key == matchString
接下来,我们将 2 个文件放入 HDFS。 /testData 是输入文件,包含制表符分隔的键值对。 /pythonAppConf.xml 是 XML 文件,我们可以在其中配置一个特定的键来匹配。
/testData
foo 1
bar 2
baz 3
/pythonAppConf.xml
<pythonAppConf>
<matchString>foo</matchString>
</pythonAppConf>
由于我们已将matchString 设置为foo,并且由于我们的输入文件仅包含一条键设置为foo 的记录,因此我们希望运行作业的输出是包含相应值的单行键入foo,即1。进行测试运行,我们确实得到了预期的结果。
> hadoop jar share/hadoop/tools/lib/hadoop-streaming-*.jar \
-D mapreduce.job.reduces=0 \
-files pythonapp \
-input /testData \
-output /streamingOut \
-mapper 'python pythonapp/Mapper.py'
> hdfs dfs -cat /streamingOut/part*
1
另一种方法是在 -files 参数中指定 HDFS 文件。这样,YARN 将在 Python 脚本启动之前将 XML 文件作为本地化资源拉取到运行容器的各个节点。然后,Python 代码可以打开 XML 文件,就好像它是工作目录中的本地文件一样。对于运行多个任务/容器的大型作业,此技术可能优于从每个任务调用 hdfs dfs -cat。
为了测试这种技术,我们可以尝试不同版本的 Process.py 模块。
进程.py
import xml.etree.ElementTree as ElementTree
pythonAppConfXmlTree = ElementTree.parse('pythonAppConf.xml')
matchString = pythonAppConfXmlTree.find('./matchString').text.strip()
def match(key):
return key == matchString
命令行调用更改为在-files 中指定一个HDFS 路径,我们再次看到了预期的结果。
> hadoop jar share/hadoop/tools/lib/hadoop-streaming-*.jar \
-D mapreduce.job.reduces=0 \
-files pythonapp,hdfs:///pythonAppConf.xml \
-input /testData \
-output /streamingOut \
-mapper 'python pythonapp/Mapper.py'
> hdfs dfs -cat /streamingOut/part*
1
Apache Hadoop 文档在这里讨论了使用 -files 选项在本地提取 HDFS 文件。
http://hadoop.apache.org/docs/r2.7.1/hadoop-streaming/HadoopStreaming.html#Working_with_Large_Files_and_Archives