首页 文章 精选 留言 我的

精选列表

搜索[并行文件系统],共10000篇文章
优秀的个人博客,低调大师

使用python的hdfs包操作分布式文件系统(HDFS)

转载请注明出处:@http://blog.csdn.net/gamer_gyt,Thinkagmer 撰写 博主微博:http://weibo.com/234654758(欢迎互撩) Github:https://github.com/thinkgamer ===================================================================================== 写在前边的话: 之前做的hadoop集群,组合了hive,hbase,sqoop,spark等开源工具,现在要对他们做一个Web的可视化操作,由于本小白只懂如何使用python做一个交互的web应用,所以这里就选择了Python的Django Django教程参考:Django从manage.py shell 到项目部署 hadoop集群操作请参考:三台PC服务器部署高可用hadoop集群 言归正传: 使用python操作hdfs本身并不难,只不过是把对应的shell 功能“翻译”成高级语言,网上大部分使用的是 pyhdfs:官方文档 hdfs:官方文档 libhdfs(比较狗血) 我这里选用的是hdfs,下边的实例都是基于hdfs包进行的 1:安装 由于我的是windows环境(linux其实也一样),只要有pip或者setup_install安装起来都是很方便的 pip install hdfs 2:Client——创建集群连接 >>> from hdfs import * >>> client = Client("http://127.0.0.1:50070") 其他参数说明: classhdfs.client.Client(url,root=None,proxy=None,timeout=None,session=None) url:ip:端口 root:制定的hdfs根目录 proxy:制定登陆的用户身份 timeout:设置的超时时间 seesion:requests.Session instance, used to emit all requests.(不是太懂,应该四用户发出请求) 这里我们着重看一下proxy这个,首先我们指定root用户连接 >>> client = Client("http://127.0.0.1:50070",root="/",timeout=100,session=False) >>> client.list("/") [u'hbase'] 看起来一切正常的样子,接下来我们指定一个别的用户,比如说gamer再看 >>> client = Client("http://127.0.0.1:50070",root="/",proxy="gamer",timeout=100,session=False) >>> client.list("/") Traceback (most recent call last): File "<stdin>", line 1, in <module> File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 893, in list statuses = self._list_status(hdfs_path).json()['FileStatuses']['FileStatus'] File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 92, in api_handler **self.kwargs File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 181, in _request return _on_error(response) File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 44, in _on_error raise HdfsError(message) hdfs.util.HdfsError: Failed to obtain user group information: org.apache.hadoop.security.authorize.AuthorizationException: User: dr.who is not allowed to impersonate gamer 这时候就抛出异常了 3:dir——查看支持的方法 >>> dir(client) ['__class__', '__delattr__', '__dict__', '__dir__', '__doc__', '__eq__', '__format__', '__ge__', '__getattribute__', '__gt__', '__hash__', '__init__', '__le__', '__lt__', '__module__', '__ne__', '__new__', '__reduce__', '__reduce_ex__', '__registry__', '__repr__', '__setattr__', '__sizeof__', '__str__', '__subclasshook__', '__weakref__', '_append', '_create', '_delete', '_get_content_summary', '_get_file_checksum', '_get_file_status', '_get_home_directory', '_list_status', '_mkdirs', '_open', '_proxy', '_rename', '_request', '_session', '_set_owner', '_set_permission', '_set_replication', '_set_times', '_timeout', 'checksum', 'content', 'delete', 'download', 'from_options', 'list', 'makedirs', 'parts', 'read', 'rename', 'resolve', 'root', 'set_owner', 'set_permission', 'set_replication', 'set_times', 'status', 'upload', 'url', 'walk', 'write'] 4:status——获取路径的具体信息 >>> client.status("/") {'accessTime': 0, 'pathSuffix': '', 'group': 'supergroup', 'type': 'DIRECTORY', 'owner': 'root', 'childrenNum': 4, 'blockSize': 0, 'fileId': 16385, 'length': 0, 'replication': 0, 'storagePolicy': 0, 'modificationTime': 1473023149031, 'permission': '777'} 其他参数:status(hdfs_path,strict=True) hdfs_path:就是hdfs路径 strict:设置为True时,如果hdfs_path路径不存在就会抛出异常,如果设置为False,如果路径为不存在,则返回None >>> client = Client("http://127.0.0.1:50070",root="/",timeout=100,session=False) >>> client.status("/gamer",strict=True) Traceback (most recent call last): File "<stdin>", line 1, in <module> File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 277, in status res = self._get_file_status(hdfs_path, strict=strict) File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 92, in api_handler **self.kwargs File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 181, in _request return _on_error(response) File "/usr/local/lib/python2.7/dist-packages/hdfs/client.py", line 44, in _on_error raise HdfsError(message) hdfs.util.HdfsError: File does not exist: /gamer >>> client.status("/gamer",strict=False) >>> 从例子中可以看出,当设置为false时,路径不存在,什么也不输出 5:list——获取指定路径的子目录信息 >>> client.list("/") ['file', 'gyt', 'hbase', 'tmp'] 其他参数:list(hdfs_path,status=False) status:为True时,也返回子目录的状态信息,默认为Flase >>> client.list("/") [u'hbase'] >>> client.list("/",status=False) [u'hbase'] >>> client.list("/",status=True) [(u'hbase', {u'group': u'supergroup', u'permission': u'755', u'blockSize': 0, u'accessTime': 0, u'pathSuffix': u'hbase', u'modificationTime': 1472986624167, u'replication': 0, u'length': 0, u'childrenNum': 7, u'owner': u'root', u'storagePolicy': 0, u'type': u'DIRECTORY', u'fileId': 16386})] >>> 6:makedirs——创建目录 >>> client.makedirs("/test") >>> client.list("/") ['file', 'gyt', 'hbase', 'test', 'tmp'] >>> client.status("/test") {'accessTime': 0, 'pathSuffix': '', 'group': 'supergroup', 'type': 'DIRECTORY', 'owner': 'dr.who', 'childrenNum': 0, 'blockSize': 0, 'fileId': 16493, 'length': 0, 'replication': 0, 'storagePolicy': 0, 'modificationTime': 1473096896947, 'permission': '755'} 其他参数:makedirs(hdfs_path,permission=None) permission:设置权限 >>> client.makedirs("/test",permission=777) >>> client.status("/test") {u'group': u'supergroup', u'permission': u'777', u'blockSize': 0, u'accessTime': 0, u'pathSuffix': u'', u'modificationTime': 1473175557340, u'replication': 0, u'length': 0, u'childrenNum': 0, u'owner': u'dr.who', u'storagePolicy': 0, u'type': u'DIRECTORY', u'fileId': 16437} 可以看出该文件夹的权限是777 7:rename—重命名 >>> client.rename("/test","/new_name") >>> client.list("/") ['file', 'gyt', 'hbase', 'new_name', 'tmp'] 格式说明:rename(hdfs_path, local_path) 8:delete—删除 >>> client.list("/") ['file', 'gyt', 'hbase', 'new_name', 'tmp'] >>> client.delete("/new_name") True >>> client.list("/") ['file', 'gyt', 'hbase', 'tmp'] 其他参数:delete(hdfs_path,recursive=False) recursive:删除文件和其子目录,设置为False如果不存在,则会抛出异常,默认为False >>> client.delete("/test",recursive=True) True >>> client.delete("/test",recursive=True) False >>> client.delete("/test") False 9:upload——上传数据 =======================分割线========================== 为什么这里需要分割线?因为在做web平台可视化操作hdfs的时候遇到了问题!错误如下: requests.exceptions.ConnectionError: HTTPConnectionPool(host='slaver1', port=50075): Max retries exceeded with url: /webhdfs/v1/thinkgamer/name.txt?op=OPEN&namenoderpcaddress=master&offset=0 (Caused by NewConnectionError ('<requests.packages.urllib3.connection.HTTPConnection object at 0x00000000043A3FD0>: Failed to establish a new connection: [Errno 11004] getaddrinfo failed',)) 对错误的理解:看其大意是Http连接太多,没有及时关闭,导致错误 (PS:网上对hdfs操作的资料比较少,大部分都只停留在基础语法层面,但对于错误的记录及解决办法少之又少) 解决办法:暂无 由于我是在windows上操作集群的,而我的集群是在服务器上部署的,所以我考虑是否在服务器上尝试下载和上传数据,果断ok >>> client.list("/") [u'hbase', u'test'] >>> client.upload("/test","/opt/bigdata/hadoop/NOTICE.txt") '/test/NOTICE.txt' >>> client.list("/") [u'hbase', u'test'] >>> client.list("/test") [u'NOTICE.txt'] 其他参数: upload ( hdfs_path , local_path , overwrite=False , n_threads=1 , temp_dir=None , chunk_size=65536,progress=None,cleanup=True,**kwargs) overwrite:是否是覆盖性上传文件 n_threads:启动的线程数目 temp_dir:当overwrite=true时,远程文件一旦存在,则会在上传完之后进行交换 chunk_size:文件上传的大小区间 progress:回调函数来跟踪进度,为每一chunk_size字节。它将传递两个参数,文件上传的路径和传输的字节数。一旦完成,-1将作为第二个参数 cleanup:如果在上传任何文件时发生错误,则删除该文件 10:download——下载 >>> client.download("/test/NOTICE.txt","/home") '/home/NOTICE.txt' >>> import os >>> os.system("ls /home") lost+found NOTICE.txt thinkgamer 0 >>> 其他参数: download ( hdfs_path , local_path , overwrite=False , n_threads=1 , temp_dir=None , **kwargs ) 参考上传 upload 11:read——读取文件 同样在windows客户端上执行依旧报错,在hadoop的节点服务器上执行 >>> with client.read("/test/NOTICE.txt") as reader: ... print reader.read() ... This product includes software developed by The Apache Software Foundation (http://www.apache.org/). >>> 其他参数: read ( *args , **kwds ) hdfs_path:hdfs路径 offset:设置开始的字节位置 length:读取的长度(字节为单位) buffer_size:用于传输数据的字节的缓冲区的大小。默认值设置在HDFS配置。 encoding:制定编码 chunk_size:如果设置为正数,上下文管理器将返回一个发生器产生的每一chunk_size字节而不是一个类似文件的对象 delimiter:如果设置,上下文管理器将返回一个发生器产生每次遇到分隔符。此参数要求指定的编码。 progress:回调函数来跟踪进度,为每一chunk_size字节(不可用,如果块大小不是指定)。它将传递两个参数,文件上传的路径和传输的字节数。称为一次与- 1作为第二个参数。 附:在对文件操作时,可能会提示错误 hdfs.util.HdfsError: Permission denied: user=dr.who, access=WRITE, inode="/test":root:supergroup:drwxr-xr-x 解决办法是:在配置文件hdfs-site.xml中加入 <property> <name>dfs.permissions</name> <value>false</value> </property> 重启集群即可 基本常用的功能也就这些了,如果需要一些特殊的功能,可以自己执行help(client.method)进行查看

优秀的个人博客,低调大师

大数据应用之HBase数据插入性能优化之多线程并行插入测试案例

一、引言: 上篇文章提起关于HBase插入性能优化设计到的五个参数,从参数配置的角度给大家提供了一个性能测试环境的实验代码。根据网友的反馈,基于单线程的模式实现的数据插入毕竟有限。通过个人实测,在我的虚拟机环境下,单线程插入数据的值约为4w/s。集群指标是:CPU双核1.83,虚拟机512M内存,集群部署单点模式。本文给出了基于多线程并发模式的,测试代码案例和实测结果,希望能给大家一些启示: 二、源程序: 1 import org.apache.hadoop.conf.Configuration; 2 import org.apache.hadoop.hbase.HBaseConfiguration; 3 import java.io.BufferedReader; 4 import java.io.File; 5 import java.io.FileNotFoundException; 6 import java.io.FileReader; 7 import java.io.IOException; 8 import java.util.ArrayList; 9 import java.util.List; 10 import java.util.Random; 11 12 import org.apache.hadoop.conf.Configuration; 13 import org.apache.hadoop.hbase.HBaseConfiguration; 14 import org.apache.hadoop.hbase.client.HBaseAdmin; 15 import org.apache.hadoop.hbase.client.HTable; 16 import org.apache.hadoop.hbase.client.HTableInterface; 17 import org.apache.hadoop.hbase.client.HTablePool; 18 import org.apache.hadoop.hbase.client.Put; 19 20 public class HBaseImportEx { 21 static Configuration hbaseConfig = null; 22 public static HTablePool pool = null; 23 public static String tableName = "T_TEST_1"; 24 static{ 25 //conf = HBaseConfiguration.create(); 26 Configuration HBASE_CONFIG = new Configuration(); 27 HBASE_CONFIG.set("hbase.master", "192.168.230.133:60000"); 28 HBASE_CONFIG.set("hbase.zookeeper.quorum", "192.168.230.133"); 29 HBASE_CONFIG.set("hbase.zookeeper.property.clientPort", "2181"); 30 hbaseConfig = HBaseConfiguration.create(HBASE_CONFIG); 31 32 pool = new HTablePool(hbaseConfig, 1000); 33 } 34 /* 35 * Insert Test single thread 36 * */ 37 public static void SingleThreadInsert()throws IOException 38 { 39 System.out.println("---------开始SingleThreadInsert测试----------"); 40 long start = System.currentTimeMillis(); 41 //HTableInterface table = null; 42 HTable table = null; 43 table = (HTable)pool.getTable(tableName); 44 table.setAutoFlush(false); 45 table.setWriteBufferSize(24*1024*1024); 46 //构造测试数据 47 List<Put> list = new ArrayList<Put>(); 48 int count = 10000; 49 byte[] buffer = new byte[350]; 50 Random rand = new Random(); 51 for(int i=0;i<count;i++) 52 { 53 Put put = new Put(String.format("row %d",i).getBytes()); 54 rand.nextBytes(buffer); 55 put.add("f1".getBytes(), null, buffer); 56 //wal=false 57 put.setWriteToWAL(false); 58 list.add(put); 59 if(i%10000 == 0) 60 { 61 table.put(list); 62 list.clear(); 63 table.flushCommits(); 64 } 65 } 66 long stop = System.currentTimeMillis(); 67 //System.out.println("WAL="+wal+",autoFlush="+autoFlush+",buffer="+writeBuffer+",count="+count); 68 69 System.out.println("插入数据:"+count+"共耗时:"+ (stop - start)*1.0/1000+"s"); 70 71 System.out.println("---------结束SingleThreadInsert测试----------"); 72 } 73 /* 74 * 多线程环境下线程插入函数 75 * 76 * */ 77 public static void InsertProcess()throws IOException 78 { 79 long start = System.currentTimeMillis(); 80 //HTableInterface table = null; 81 HTable table = null; 82 table = (HTable)pool.getTable(tableName); 83 table.setAutoFlush(false); 84 table.setWriteBufferSize(24*1024*1024); 85 //构造测试数据 86 List<Put> list = new ArrayList<Put>(); 87 int count = 10000; 88 byte[] buffer = new byte[256]; 89 Random rand = new Random(); 90 for(int i=0;i<count;i++) 91 { 92 Put put = new Put(String.format("row %d",i).getBytes()); 93 rand.nextBytes(buffer); 94 put.add("f1".getBytes(), null, buffer); 95 //wal=false 96 put.setWriteToWAL(false); 97 list.add(put); 98 if(i%10000 == 0) 99 { 100 table.put(list); 101 list.clear(); 102 table.flushCommits(); 103 } 104 } 105 long stop = System.currentTimeMillis(); 106 //System.out.println("WAL="+wal+",autoFlush="+autoFlush+",buffer="+writeBuffer+",count="+count); 107 108 System.out.println("线程:"+Thread.currentThread().getId()+"插入数据:"+count+"共耗时:"+ (stop - start)*1.0/1000+"s"); 109 } 110 111 112 /* 113 * Mutil thread insert test 114 * */ 115 public static void MultThreadInsert() throws InterruptedException 116 { 117 System.out.println("---------开始MultThreadInsert测试----------"); 118 long start = System.currentTimeMillis(); 119 int threadNumber = 10; 120 Thread[] threads=new Thread[threadNumber]; 121 for(int i=0;i<threads.length;i++) 122 { 123 threads[i]= new ImportThread(); 124 threads[i].start(); 125 } 126 for(int j=0;j< threads.length;j++) 127 { 128 (threads[j]).join(); 129 } 130 long stop = System.currentTimeMillis(); 131 132 System.out.println("MultThreadInsert:"+threadNumber*10000+"共耗时:"+ (stop - start)*1.0/1000+"s"); 133 System.out.println("---------结束MultThreadInsert测试----------"); 134 } 135 136 /** 137 * @param args 138 */ 139 public static void main(String[] args) throws Exception{ 140 // TODO Auto-generated method stub 141 //SingleThreadInsert(); 142 MultThreadInsert(); 143 144 145 } 146 147 public static class ImportThread extends Thread{ 148 public void HandleThread() 149 { 150 //this.TableName = "T_TEST_1"; 151 152 153 } 154 // 155 public void run(){ 156 try{ 157 InsertProcess(); 158 } 159 catch(IOException e){ 160 e.printStackTrace(); 161 }finally{ 162 System.gc(); 163 } 164 } 165 } 166 167 } 三、说明 1.线程数设置需要根据本集群硬件参数,实际测试得出。否则线程过多的情况下,总耗时反而是下降的。 2.单笔提交数对性能的影响非常明显,需要在自己的环境下,找到最理想的数值,这个需要与单条记录的字节数相关。 四、测试结果 ---------开始MultThreadInsert测试---------- 线程:8插入数据:10000共耗时:1.328s线程:16插入数据:10000共耗时:1.562s线程:11插入数据:10000共耗时:1.562s线程:10插入数据:10000共耗时:1.812s线程:13插入数据:10000共耗时:2.0s线程:17插入数据:10000共耗时:2.14s线程:14插入数据:10000共耗时:2.265s线程:9插入数据:10000共耗时:2.468s线程:15插入数据:10000共耗时:2.562s线程:12插入数据:10000共耗时:2.671sMultThreadInsert:100000共耗时:2.703s---------结束MultThreadInsert测试---------- 备注:该技术专题讨论正在群Hadoop高级交流群:293503507同步直播中,敬请关注。 作者:张子良 出处:http://www.cnblogs.com/hadoopdev 本文版权归作者所有,欢迎转载,但未经作者同意必须保留此段声明,且在文章页面明显位置给出原文连接,否则保留追究法律责任的权利。

资源下载

更多资源
Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册