【问题标题】:Problem in creating a separate function to read file in Hadoop Streaming在 Hadoop Streaming 中创建单独的函数来读取文件时出现问题
【发布时间】:2011-10-09 07:04:56
【问题描述】:

我在创建一个单独的函数来读取 Hadoop Streaming 中的文件时遇到问题。

mapper.py:Works well (very inefficient)

#!/usr/bin/env python 导入系统 定义主(): 对于 sys.stdin 中的行: line = line.strip() # 每行只包含一个单词,5+百万行 filename = "my_dict.txt" # 包含 7+ 百万字 f = 打开(文件名,“r”) 对于 f 中的第 1 行: line1 = line1.strip() 如果第 1 行 == 行: 打印 '%s\t%s' % (line1, 1) 如果 __name__ == '__main__': 主要的()

我的文件读取功能

mapper.py:not working

#!/usr/bin/env python 导入系统 def read_dict(): 列表D=[] 文件名 = “my_dict.txt” f = 打开(文件名,“r”) 对于 f 中的第 1 行: listD.append(line1.strip()) 返回列表D 定义主(): listDM = set(read_dict()) 对于 sys.stdin 中的行: line = line.strip() 对于 listDM 中的第 1 行: 如果第 1 行 == 行: 打印 '%s\t%s' % (line1, 1) 如果 __name__ == '__main__': 主要的()

我的错误日志:

hadoop_admin@gml-VirtualBox:/usr/local/hadoop$ sh myScripts/runHadoopStream.sh 删除 hdfs://localhost:9000/user/hadoop_admin/output packageJobJar: [/var/www/HMS/my_dict.txt, /usr/local/hadoop/mapper.py, /usr/local/hadoop/reducer.py, /tmp/hadoop-hadoop_admin/hadoop-unjar6634198925772314314/] [] /tmp/streamjob1880303974118879660.jar tmpDir=null 20 年 11 月 7 日 18:40:08 信息 mapred.FileInputFormat:要处理的总输入路径:2 20 年 11 月 7 日 18:40:08 信息流。StreamJob:getLocalDirs():[/tmp/hadoop-hadoop_admin/mapred/local] 20 年 11 月 7 日 18:40:08 信息流。StreamJob:正在运行的作业:job_201107181559_0091 11/07/20 18:40:08 INFO streaming.StreamJob:要终止此作业,请运行: 11/07/20 18:40:08 INFO streaming.StreamJob: /usr/local/hadoop/bin/../bin/hadoop job -Dmapred.job.tracker=localhost:9001 -kill job_201107181559_0091 11/07/20 18:40:08 INFO streaming.StreamJob:跟踪 URL:http://localhost:50030/jobdetails.jsp?jobid=job_201107181559_0091 20 年 11 月 7 日 18:40:09 信息流。StreamJob:地图 0% 减少 0% 20 年 11 月 7 日 18:40:41 信息流。StreamJob:地图 1% 减少 0% 20 年 11 月 7 日 18:40:47 信息流。StreamJob:地图 0% 减少 0% 20 年 11 月 7 日 18:41:05 信息流。StreamJob:地图 1% 减少 0% 20 年 11 月 7 日 18:41:08 信息流。StreamJob:地图 0% 减少 0% 20 年 11 月 7 日 18:41:26 信息流。StreamJob:地图 1% 减少 0% 20 年 11 月 7 日 18:41:29 信息流。StreamJob:地图 0% 减少 0% 20 年 11 月 7 日 18:41:48 信息流。StreamJob:地图 1% 减少 0% 20 年 11 月 7 日 18:41:51 信息流。StreamJob:地图 0% 减少 0% 20 年 11 月 7 日 18:41:57 信息流。StreamJob:地图 100% 减少 100% 11/07/20 18:41:57 INFO streaming.StreamJob:要终止此作业,请运行: 11/07/20 18:41:57 INFO streaming.StreamJob: /usr/local/hadoop/bin/../bin/hadoop job -Dmapred.job.tracker=localhost:9001 -kill job_201107181559_0091 11/07/20 18:41:57 INFO streaming.StreamJob:跟踪 URL:http://localhost:50030/jobdetails.jsp?jobid=job_201107181559_0091 20 年 11 月 7 日 18:41:57 错误流。StreamJob:作业不成功! 20 年 11 月 7 日 18:41:57 信息流。StreamJob:killJob... 流式传输作业失败!

用于运行 Hadoop Streaming 的 Shell 脚本:

bin/hadoop dfs -rmr 输出 bin/hadoop jar contrib/streaming/hadoop-*-streaming.jar -file /var/www/HMS/my_dict.txt -file /usr/local/hadoop/mapper.py -mapper /usr/local/hadoop/mapper. py -file /usr/local/hadoop/reducer.py -reducer /usr/local/hadoop/reducer.py -input 输入/ -output 输出/

【问题讨论】:

  • 想知道我的脚本是否符合您的要求?
  • 是的,非常感谢 agf :)

标签: python hadoop


【解决方案1】:

试试这个简化的脚本:

#!/usr/bin/env python

import sys

def main():
    filename = "my_dict.txt"
    listfile = open(filename)
    # doesn't create an itermediate list
    listDM = set(line.strip() for line in listfile)
    # less Pythonic but significantly faster
    # still doesn't create an intermediate list
    # listDM = set(imap(str.strip, listfile))
    listfile.close()  
    for line in sys.stdin:
        line = line.strip()
        if line in listDM:
            print  '%s\t%d' % (line, 1) 

if __name__ == '__main__':
    main()

如果您使用更快的注释掉替代方案,您需要from itertools import imap

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-02-22
    • 2012-12-26
    • 1970-01-01
    • 2010-12-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多