当前位置:首页 > PHP教程 > PHP总结归纳

地图reduce与mysql交互

mapreduce 与mysql 交互 mapreduce与mysql交互 分类:?hadoop之旅 2012-08-29 16:12 ? 377人阅读 ? 评论(0) ? 收藏 ? 举报 mapreducemysql数据库hadoopstringjdbc? 目录(?) [+] ? mapreduce技术推出后,曾遭到关系数据库研究者的挑剔和批评,认为mapreduce不

mapreduce 与mysql 交互

  mapreduce技术推出后,曾遭到关系数据库研究者的挑剔和批评,认为mapreduce不具备有类似于关系数据库中的结构化数据存储和处理能力。为此,google和mapreduce社区进行了很多努力。一方面,他们设计了类似于关系数据中结构化数据表的技术(google的bigtable,hadoop的hbase)提供一些粗粒度的结构化数据存储和处理能力;另一方面,为了增强与关系数据库的集成能力,hadoop mapreduce提供了相应的访问关系数据库库的编程接口。

  mapreduce与mysql交互的整体架构如下图所示。

?

图2-1整个环境的架构

  具体到mapreduce框架读/写数据库,有2个主要的程序分别是?dbinputformat和dboutputformat,dbinputformat 对应的是sql语句select,而dboutputformat 对应的是?inster/update,使用dbinputformat和dboutputforma时候需要实现inputformat这个抽象类,这个抽象类含有getsplits()和createrecordreader()抽象方法,在dbinputformat类中由 protected string getcountquery() 方法传入结果集的个数,getsplits()方法再确定输入的切分原则,利用sql中的 limit 和 offset 进行切分获得数据集的范围 ,请参考dbinputformat源码中public inputsplit[] getsplits(jobconf job, int chunks) throws ioexception的方法,在dbinputformat源码中createrecordreader()则可以按一定格式读取相应数据。

??????1)建立关系数据库连接

  • dbconfiguration:提供数据库配置和创建连接的接口。

????? dbconfiguration类中提供了一个静态方法创建数据库连接:

?

public static void configuredb(job?job,string?driverclass,string?dburl,string?username,string?password)

?

????? 其中,job为当前准备执行的作业,driverclasss为数据库厂商提供的访问其数据库的驱动程序,dburl为运行数据库的主机的地址,username和password分别为数据库提供访问地用户名和相应的访问密码。

??????2)相应的从关系数据库查询和读取数据的接口

  • dbinputformat:提供从数据库读取数据的格式。
  • dbrecordreader:提供读取数据记录的接口。

  3)相应的向关系数据库直接输出结果的编程接口

  • dboutputformat:提供向数据库输出数据的格式。
  • dbrecordwrite:提供数据库写入数据记录的接口。

  数据库连接完成后,即可完成从mapreduce程序向关系数据库写入数据的操作。为了告知数据库将写入哪个表中的哪些字段,dboutputformat中提供了一个静态方法来指定需要写入的数据表和字段:

?

public static void setoutput(job job,string tablename,string ... fieldname)

?

????? 其中,tablename指定即将写入的数据表,后续参数将指定哪些字段数据将写入该表。

1.1 从数据库中输入数据

??????虽然hadoop允许从数据库中直接读取数据记录作为mapreduce的输入,但处理效率较低,而且大量频繁地从mapreduce程序中查询和读取关系数据库可能会大大增加数据库的访问负载,因此dbinputformat仅适合读取小量数据记录的计算和应用,不适合数据仓库联机数据分析大量数据的读取处理。

????? 读取大量数据记录一个更好的解决办法是:用数据库中的dump工具将大量待分析数据输出为文本数据文件,并上载到hdfs中进行处理。

?

??????1)首先创建要读入的数据

  • windows环境

  首先创建数据库"school",使用下面命令进行:

?

create database school;

?

????? 然后通过以下几句话,把我们事先准备好的sql语句(student.sql事先放到了d盘目录)导入到刚创建的"school"数据库中。用到的命令如下:

?

use school;

source d:student.sql

?

????? "student.sql"中的内容如下所示:

?

drop table if exists `school`.`student`;

?

create table `school`.`student` (

`id` int(11) not null default '0',

`name` varchar(20) default null,

`sex` varchar(10) default null,

`age` int(10) default null,

primary key (`id`)

) engine=innodb default charset=utf8;

?

insert into `student` values ('201201', '张三', '男', '21');

insert into `student` values ('201202', '李四', '男', '22');

insert into `student` values ('201203', '王五', '女', '20');

insert into `student` values ('201204', '赵六', '男', '21');

insert into `student` values ('201205', '小红', '女', '19');

insert into `student` values ('201206', '小明', '男', '22');

?

????? 执行结果如下所示:

?

?

????? 查询刚才创建的数据库表"student"的内容。

?

?

????? 结果发现显示是乱码,记得我当时是设置的utf-8,怎么就出现乱码了呢?其实我们使用的操作系统的系统为中文,且它的默认编码是gbk,而mysql的编码有两种,它们分别是:

  【client】:客户端的字符集。客户端默认字符集。当客户端向服务器发送请求时,请求以该字符集进行编码。

  【mysqld】:服务器字符集,默认情况下所采用的。

?

????? 找到安装mysql目录,比如我们的安装目录为:

?

e:hadoopworkplatmysql server 5.5

?

????? 从中找到"my.ini"配置文件,最终发现my.ini里的2个character_set把client改成gbk,把server改成utf8就可以了。

??? 【client】端:

?

[client]

port=3306

[mysql]

default-character-set=gbk

?

??? 【mysqld】端:

?

[mysqld]

# the default character set that will be used when a new schema or table is

# created and no character set is defined

character-set-server=utf8

?

????? 按照上面修改完之后,重启mysql服务。

?

?

????? 此时在windows下面的数据库表已经准备完成了。

?

  • linux环境

  首先通过"flashfxp"把我们刚才的"student.sql"上传到"/home/hadoop"目录下面,然后按照上面的语句创建"school"数据库。

?

  

????? 查看我们上传的"student.sql"内容:

?

  

??? ? 创建"school"数据库,并导入"student.sql"语句。

?

  

?

????? 显示数据库"school"中的表"student"信息。

?

  

??? ?显示表"student"中的内容。

?

  

?

????? 到此为止在"windows"和"linux"两种环境下面都创建了表"student"表,并初始化了值。下面就开始通过mapreduce读取mysql库中表"student"的信息。

??????2)使mysql能远程连接

????? mysql默认是允许别的机器进行远程访问地,为了使hadoop集群能访问mysql数据库,所以进行下面操作。

  • 用mysql用户"root"登录。

?

mysql -u root -p

?

  • 使用下面语句进行授权,赋予任何主机访问数据的权限。

?

grant all privileges on *.* to 'root'@'%' identified by 'hadoop' with grant option;

?

  • 刷新,使之立即生效。

?

flush privileges;

?

????? 执行结果如下图。

????? windows下面:

?

  

????? linux下面:

?

  

???? ?到目前为止,如果连接win7上面的mysql数据库还不行,大家还应该记得前面在linux下面关掉了防火墙,但是我们在win7下对防火墙并没有做任何处理,如果不对防火墙做处理,即使执行了上面的远程授权,仍然不能连接。下面是设置win7上面的防火墙,使远程机器能通过3306端口访问mysql数据库。

??????解决方法:只要在'入站规则'上建立一个3306端口即可。

  执行顺序:控制面板à管理工具à高级安全的windows防火墙à入站规则

  然后新建规则à选择'端口'à在'特定本地端口'上输入一个'3306'?à选择'允许连接'=>选择'域'、'专用'、'公用'=>给个名称,如:mysqlinput

?

??????3)对jdbc的jar包处理

???? ?因为程序虽然用eclipse编译运行但最终要提交到hadoop集群上,所以jdbc的jar必须放到hadoop集群中。有两种方式:

????? (1)在每个节点下的${hadoop_home}/lib下添加该包,重启集群,一般是比较原始的方法。

????? 我们的hadoop安装包在"/usr/hadoop",所以把jar放到"/usr/hadoop/lib"下面,然后重启,记得是hadoop集群中所有的节点都要放,因为执行分布式是程序是在每个节点本地机器上进行。

???? ?(2)在hadoop集群的分布式文件系统中创建"/lib"文件夹,并把我们的的jdbc的jar包上传上去,然后在主程序添加如下语句,就能保证hadoop集群中所有的节点都能使用这个jar包。因为这个jar包放在了hdfs上,而不是本地系统,这个要理解清楚。

?

distributedcache.addfiletoclasspath(new path("/lib/mysql-connector-java-5.1.18-bin.jar"), conf);

?

?? ?? 我们用的jdbc的jar如下所示:

?

mysql-connector-java-5.1.18-bin.jar

?

???? ?通过eclipse下面的dfs locations进行创建"/lib"文件夹,并上传jdbc的jar包。执行结果如下:

  

??????备注:我们这里采用了第二种方式。

???? ?4)源程序代码如下所示

?

package?com.hebut.mr;

?

import?java.io.ioexception;

import?java.io.datainput;

import?java.io.dataoutput;

import?java.sql.connection;

import?java.sql.drivermanager;

import?java.sql.preparedstatement;

import?java.sql.resultset;

import?java.sql.sqlexception;

?

import?org.apache.hadoop.filecache.distributedcache;

import?org.apache.hadoop.fs.path;

import?org.apache.hadoop.io.longwritable;

import?org.apache.hadoop.io.text;

import?org.apache.hadoop.io.writable;

import?org.apache.hadoop.mapred.jobclient;

import?org.apache.hadoop.mapred.jobconf;

import?org.apache.hadoop.mapred.mapreducebase;

import?org.apache.hadoop.mapred.mapper;

import?org.apache.hadoop.mapred.outputcollector;

import?org.apache.hadoop.mapred.fileoutputformat;

import?org.apache.hadoop.mapred.reporter;

import?org.apache.hadoop.mapred.lib.identityreducer;

import?org.apache.hadoop.mapred.lib.db.dbwritable;

import?org.apache.hadoop.mapred.lib.db.dbinputformat;

import?org.apache.hadoop.mapred.lib.db.dbconfiguration;

?

public?class?readdb {

?

????public?static?class?map?extends?mapreducebase?implements

??????????? mapper {

?

????????//?实现map函数

????????public?void?map(longwritable key, studentrecord value,

??????? outputcollector collector, reporter reporter)

????????????????throws?ioexception {

??????????? collector.collect(new?longwritable(value.id),

????????????????????new?text(value.tostring()));

??????? }

?

??? }

?

????public?static?class?studentrecord?implements?writable, dbwritable {

????????public?int?id;

????????public?string?name;

????????public?string?sex;

????????public?int?age;

?

????????@override

????????public?void?readfields(datainput in)?throws?ioexception {

????????????this.id?= in.readint();

????????????this.name?= text.readstring(in);

????????????this.sex?= text.readstring(in);

????????????this.age?= in.readint();

??????? }

?

????????@override

????????public?void?write(dataoutput out)?throws?ioexception {

??????????? out.writeint(this.id);

??????????? text.writestring(out,?this.name);

??????????? text.writestring(out,?this.sex);

??????????? out.writeint(this.age);

??????? }

?

????????@override

????????public?void?readfields(resultset result)?throws?sqlexception {

????????????this.id?= result.getint(1);

????????????this.name?= result.getstring(2);

????????????this.sex?= result.getstring(3);

????????????this.age?= result.getint(4);

??????? }

?

????????@override

????????public?void?write(preparedstatement stmt)?throws?sqlexception {

??????????? stmt.setint(1,?this.id);

??????????? stmt.setstring(2,?this.name);

??????????? stmt.setstring(3,?this.sex);

??????????? stmt.setint(4,?this.age);

??????? }

?

????????@override

????????public?string tostring() {

????????????return?new?string("学号:"?+?this.id?+?"_姓名:"?+?this.name

??????????????????? +?"_性别:"+?this.sex?+?"_年龄:"?+?this.age);

??????? }

??? }

?

????public?static?void?main(string[] args)?throws?exception {

?

??????? jobconf conf =?new?jobconf(readdb.class);

?

????????//?这句话很关键

??????? conf.set("mapred.job.tracker",?"192.168.1.2:9001");

?

??????? //?非常重要,值得关注

??????? distributedcache.addfiletoclasspath(new?path(

?????????"/lib/mysql-connector-java-5.1.18-bin.jar"), conf);

?

????????//?设置输入类型

??????? conf.setinputformat(dbinputformat.class);

?

????????//?设置输出类型

??????? conf.setoutputkeyclass(longwritable.class);

??????? conf.setoutputvalueclass(text.class);

?

????????//?设置map和reduce类

??????? conf.setmapperclass(map.class);

??????? conf.setreducerclass(identityreducer.class);

?

????????//?设置输出目录

??????? fileoutputformat.setoutputpath(conf,?new?path("rdb_out"));

?

????????//?建立数据库连接

??????? dbconfiguration.configuredb(conf,?"com.mysql.jdbc.driver",

????????????"jdbc:mysql://192.168.1.24:3306/school",?"root",?"hadoop");

?

????????//?读取"student"表中的数据

??????? string[] fields = {?"id",?"name",?"sex",?"age"?};

??????? dbinputformat.setinput(conf, studentrecord.class,?"student",?null,"id", fields);

?

??????? jobclient.runjob(conf);

??? }

}

?

???? ?备注:由于hadoop1.0.0新的api对关系型数据库暂不支持,只能用旧的api进行,所以下面的"向数据库中输出数据"也是如此。

?

???? ?5)运行结果如下所示

???? ?经过上面的设置后,已经通过连接win7和linux上的mysql数据库,执行结果都一样。唯独变得就是代码中"dbconfiguration.configuredb"中mysql数据库所在机器的ip地址。

?

?

1.2 向数据库中输出数据

???? ?基于数据仓库的数据分析和挖掘输出结果的数据量一般不会太大,因而可能适合于直接向数据库写入。我们这里尝试与"wordcount"程序相结合,把单词统计的结果存入到关系型数据库中。

???? ?1)创建写入的数据库表

??? ? 我们还使用刚才创建的数据库"school",只是在里添加一个新的表"wordcount",还是使用下面语句执行:

?

use school;

source?sql脚本全路径

?

??? ? 下面是要创建的"wordcount"表的sql脚本。

?

drop table if exists `school`.`wordcount`;

?

create table `school`.`wordcount` (

`id` int(11) not null auto_increment,

`word` varchar(20) default null,

`number` int(11) default null,

primary key (`id`)

) engine=innodb default charset=utf8;

?

???? ?执行效果如下所示:

  • windows环境

  • linux环境

?

??????2)程序源代码如下所示

?

package?com.hebut.mr;

?

import?java.io.ioexception;

import?java.io.datainput;

import?java.io.dataoutput;

import?java.sql.preparedstatement;

import?java.sql.resultset;

import?java.sql.sqlexception;

import?java.util.iterator;

import?java.util.stringtokenizer;

?

import?org.apache.hadoop.filecache.distributedcache;

import?org.apache.hadoop.fs.path;

import?org.apache.hadoop.io.intwritable;

import?org.apache.hadoop.io.text;

import?org.apache.hadoop.io.writable;

import?org.apache.hadoop.mapred.fileinputformat;

import?org.apache.hadoop.mapred.jobclient;

import?org.apache.hadoop.mapred.jobconf;

import?org.apache.hadoop.mapred.mapreducebase;

import?org.apache.hadoop.mapred.mapper;

import?org.apache.hadoop.mapred.outputcollector;

import?org.apache.hadoop.mapred.reducer;

import?org.apache.hadoop.mapred.reporter;

import?org.apache.hadoop.mapred.textinputformat;

import?org.apache.hadoop.mapred.lib.db.dboutputformat;

import?org.apache.hadoop.mapred.lib.db.dbwritable;

import?org.apache.hadoop.mapred.lib.db.dbconfiguration;

?

public?class?writedb {

????// map处理过程

????public?static?class?map?extends?mapreducebase?implements

??????????? mapper {

?

????????private?final?static?intwritable?one?=?new?intwritable(1);

????????private?text?word?=?new?text();

?

????????@override

????????public?void?map(object key, text value,

??????????? outputcollector output, reporter reporter)

????????????????throws?ioexception {

??????????? string line = value.tostring();

??????????? stringtokenizer tokenizer =?new?stringtokenizer(line);

????????????while?(tokenizer.hasmoretokens()) {

????????????????word.set(tokenizer.nexttoken());

??????????????? output.collect(word,?one);

??????????? }

??????? }

??? }

?

????// combine处理过程

????public?static?class?combine?extends?mapreducebase?implements

??????????? reducer {

?

????????@override

????????public?void?reduce(text key, iterator values,

??????????? outputcollector output, reporter reporter)

????????????????throws?ioexception {

????????????int?sum = 0;

????????????while?(values.hasnext()) {

??????????????? sum += values.next().get();

??????????? }

??????????? output.collect(key,?new?intwritable(sum));

??????? }

??? }

?

????// reduce处理过程

????public?static?class?reduce?extends?mapreducebase?implements

??????????? reducer {

?

????????@override

????????public?void?reduce(text key, iterator values,

??????????? outputcollector collector, reporter reporter)

????????????????throws?ioexception {

?

????????????int?sum = 0;

????????????while?(values.hasnext()) {

??????????????? sum += values.next().get();

??????????? }

?

??????????? wordrecord wordcount =?new?wordrecord();

??????????? wordcount.word?= key.tostring();

??????????? wordcount.number?= sum;

?

??????????? collector.collect(wordcount,?new?text());

??????? }

??? }

?

????public?static?class?wordrecord?implements?writable, dbwritable {

????????public?string?word;

????????public?int?number;

?

????????@override

????????public?void?readfields(datainput in)?throws?ioexception {

????????????this.word?= text.readstring(in);

????????????this.number?= in.readint();

??????? }

?

????????@override

????????public?void?write(dataoutput out)?throws?ioexception {

??????????? text.writestring(out,?this.word);

??????????? out.writeint(this.number);

??????? }

?

????????@override

????????public?void?readfields(resultset result)?throws?sqlexception {

????????????this.word?= result.getstring(1);

????????????this.number?= result.getint(2);

??????? }

?

????????@override

????????public?void?write(preparedstatement stmt)?throws?sqlexception {

??????????? stmt.setstring(1,?this.word);

??????????? stmt.setint(2,?this.number);

??????? }

??? }

?

????public?static?void?main(string[] args)?throws?exception {

?

??????? jobconf conf =?new?jobconf(writedb.class);

?

????????//?这句话很关键

??????? conf.set("mapred.job.tracker",?"192.168.1.2:9001");

?

??????? distributedcache.addfiletoclasspath(new?path(

????????????????"/lib/mysql-connector-java-5.1.18-bin.jar"), conf);

?

????????//?设置输入输出类型

??????? conf.setinputformat(textinputformat.class);

??????? conf.setoutputformat(dboutputformat.class);

????????//?不加这两句,通不过,但是网上给的例子没有这两句。

??????? conf.setoutputkeyclass(text.class);

??????? conf.setoutputvalueclass(intwritable.class);

?

????????//?设置map和reduce类

??????? conf.setmapperclass(map.class);

??????? conf.setcombinerclass(combine.class);

??????? conf.setreducerclass(reduce.class);

?

????????//?设置输如目录

??????? fileinputformat.setinputpaths(conf,?new?path("wdb_in"));

?

????????//?建立数据库连接

??????? dbconfiguration.configuredb(conf,?"com.mysql.jdbc.driver",

????????????"jdbc:mysql://192.168.1.24:3306/school",?"root",?"hadoop");

?

????????//?写入"wordcount"表中的数据

??????? string[] fields = {?"word",?"number"?};

??????? dboutputformat.setoutput(conf,?"wordcount", fields);

?

??????? jobclient.runjob(conf);

??? }

}

?

???? ?3)运行结果如下所示

  • windows环境

  测试数据:

(1)file1.txt

?

hello word

hello hadoop

?

????(2)file2.txt

?

虾皮 hadoop

虾皮 word

软件 软件

?

???? ?运行结果:

?

???? ?我们发现上图中出现了"?",后来查找原来是因为我的测试数据时在windows用记事本写的然后保存为"utf-8",在保存时为了区分编码,自动在前面加了一个"bom",但是不会显示任何结果。然而我们的代码把它识别为"?"进行处理。这就出现了上面的结果,如果我们在每个要处理的文件前面的第一行加一个空格,结果就成如下显示:

?

?? ?? 接着又做了一个测试,在linux上面用下面命令创建了一个文件,并写上中文内容。结果显示并没有出现"?",而且网上说不同的记事本软件(emeditor、ue)保存为"utf-8"就没有这个问题。经过修改之后的map类,就能够正常识别了。

?

????// map处理过程

????public?static?class?map?extends?mapreducebase?implements

??????????? mapper {

?

????????private?final?static?intwritable?one?=?new?intwritable(1);

????????private?text?word?=?new?text();

?

????????@override

????????public?void?map(object key, text value,

??????????? outputcollector output, reporter reporter)

????????????????throws?ioexception {

??????????? string line = value.tostring();

???????????

????????????//处理记事本utf-8的bom问题

????????????if?(line.getbytes().length?> 0) {

????????????????if?((int) line.charat(0) == 65279) {

??????????????????? line = line.substring(1);

??????????????? }

??????????? }

???????????

??????????? stringtokenizer tokenizer =?new?stringtokenizer(line);

????????????while?(tokenizer.hasmoretokens()) {

????????????????word.set(tokenizer.nexttoken());

??????????????? output.collect(word,?one);

??????????? }

??????? }

??? }

?

???? ?处理之后的结果:

?

?

???? ?从上图中得知,我们的问题已经解决了,因此,在编辑、更改任何文本文件时,请务必使用不会乱加bom的编辑器。linux下的编辑器应该都没有这个问题。windows下,请勿使用记事本等编辑器。推荐的编辑器是: editplus 2.12版本以上; emeditor; ultraedit(需要取消'添加bom'的相关选项); dreamweaver(需要取消'添加bom'的相关选项) 等。

  对于已经添加了bom的文件,要取消的话,可以用以上编辑器另存一次。(editplus需要先另存为gb,再另存为utf-8。) dw解决办法如下: 用dw打开指定文件,按ctrl+jà标题/编码à编码选择"utf-8",去掉"包括unicode签名(bom)"勾选à保存/另存为,即可。

??? 国外有一个牛人已经把这个问题解决了,使用"unicodeinputstream"、"unicodereader"。

??? 地址:http://koti.mbnet.fi/akini/java/unicodereader/

??? 示例:java读带有bom的utf-8文件乱码原因及解决方法

??? 代码:http://download.csdn.net/detail/xia520pi/4146123

?

  • linux环境

  测试数据:

????(1)file1.txt

?

mapreduce is simple

?

????(2)file2.txt

?

mapreduce is powerful is simple

?

????(3)file2.txt

?

hello mapreduce bye mapreduce

?

??? ? 运行结果:

?

?

?? ?? 到目前为止,mapreduce与关系型数据库交互已经结束,从结果中得知,目前新版的api还不能很好的支持关系型数据库的操作,上面两个例子都是使用的旧版的api。关于更多的mysql操作。?

??????终于完成,期间遇到的关键问题如下:

?

  • mysql的jdbc的jar存放问题。
  • win7对mysql防火墙的设置。
  • linux中mysql变更目录不能启动。
  • mapreduce处理带bom的utf-8问题。
  • 设置mysql可以远程访问。
  • mysql处理中文乱码问题。

?

  从这几天对mapreduce的了解,发现其实hadoop对关系型数据库的处理还不是很强,主要是hadoop和关系型数据做的事不是同一类型,各有所特长。下面几期我们将对hadoop里的hbase和hive进行全面了解。

.syntaxhighlighter{padding-top:20px;padding-bottom:20px;}

【说明】本文章由站长整理发布,文章内容不代表本站观点,如文中有侵权行为,请与本站客服联系(QQ:)!