首页 文章 精选 留言 我的

精选列表

搜索[流计算],共10000篇文章
优秀的个人博客,低调大师

Mysql 流增量写入 Hdfs(二) --Storm + hdfs 的流式处理

一. 概述 上一篇我们介绍了如何将数据从 mysql 抛到 kafka,这次我们就专注于利用 storm 将数据写入到 hdfs 的过程,由于 storm 写入 hdfs 的可定制东西有些多,我们先不从 kafka 读取,而先自己定义一个 Spout 数据充当数据源,下章再进行整合。这里默认你是拥有一定的 storm 知识的基础,起码知道 Spout 和 bolt 是什么。 写入 hdfs 可以有以下的定制策略: 自定义写入文件的名字 定义写入内容格式 满足给定条件后更改写入的文件 更改写入文件时触发的 Action 本篇会先说明如何用 storm 写入 HDFS,写入过程一些 API 的描述,以及最后给定一个例子: storm 每接收到 10 个 Tuple 后就会改变 hdfs 写入文件,新文件的名字就是第几次改变。 ps:storm 版本:1.1.1 。Hadoop 版本:2.7.4 。 接下来我们首先看看 Storm 如何写入 HDFS 。 二. Storm 写入 HDFS Storm 官方有提供了相应的 API 让我们可以使用。可以通过创建 HdfsBolt 以及定义相应的规则,即可写入 HDFS 。 首先通过 maven 配置依赖以及插件。 <properties> <storm.version>1.1.1</storm.version> </properties> <dependencies> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-core</artifactId> <version>${storm.version}</version> <!--<scope>provided</scope>--> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>log4j-over-slf4j</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>commons-collections</groupId> <artifactId>commons-collections</artifactId> <version>3.2.1</version> </dependency> <dependency> <groupId>com.google.guava</groupId> <artifactId>guava</artifactId> <version>15.0</version> </dependency> <!--hadoop模块--> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>2.7.4</version> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-hdfs</artifactId> <version>2.7.4</version> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> </exclusion> </exclusions> </dependency> <!-- https://mvnrepository.com/artifact/org.apache.storm/storm-hdfs --> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-hdfs</artifactId> <version>1.1.1</version> <!--<scope>test</scope>--> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.5.1</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> <plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>1.2.1</version> <executions> <execution> <goals> <goal>exec</goal> </goals> </execution> </executions> <configuration> <executable>java</executable> <includeProjectDependencies>true</includeProjectDependencies> <includePluginDependencies>false</includePluginDependencies> <classpathScope>compile</classpathScope> <mainClass>com.learningstorm.kafka.KafkaTopology</mainClass> </configuration> </plugin> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>1.7</version> <configuration> <createDependencyReducedPom>true</createDependencyReducedPom> </configuration> <executions> <execution> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass></mainClass> </transformer> </transformers> </configuration> </execution> </executions> </plugin> </plugins> </build> 这里要提一下,如果要打包部署到集群上的话,打包的插件需要使用 maven-shade-plugin 这个插件,然后使用 maven Lifecycle 中的 package 打包。而不是用 Maven-assembly-plugin 插件进行打包。 因为使用 Maven-assembly-plugin 的时候,会将所有依赖的包unpack,然后在pack,这样就会出现,同样的文件被覆盖的情况。发布到集群上的时候就会报 No FileSystem for scheme: hdfs 的错 。 然后是使用 HdfsBolt 写入 Hdfs。这里来看看官方文档中的例子吧。 // 使用 "|" 来替代 ",",来进行字符分割 RecordFormat format = new DelimitedRecordFormat() .withFieldDelimiter("|"); // 每输入 1k 后将内容同步到 Hdfs 中 SyncPolicy syncPolicy = new CountSyncPolicy(1000); // 当文件大小达到 5MB ,转换写入文件,即写入到一个新的文件中 FileRotationPolicy rotationPolicy = new FileSizeRotationPolicy(5.0f, Units.MB); //当转换写入文件时,生成新文件的名字并使用 FileNameFormat fileNameFormat = new DefaultFileNameFormat() .withPath("/foo/"); HdfsBolt bolt = new HdfsBolt() .withFsUrl("hdfs://localhost:9000") .withFileNameFormat(fileNameFormat) .withRecordFormat(format) .withRotationPolicy(rotationPolicy) .withSyncPolicy(syncPolicy); //生成该 bolt topologyBuilder.setBolt("hdfsBolt", bolt, 5).globalGrouping("randomStrSpout"); 到这里就结束了。可以将 HdfsBolt 当作一个 Storm 中特殊一些的 bolt 即可。这个 bolt 的作用即使根据接收信息写入 Hdfs。 而在新建 HdfsBolt 中,Storm 为我们提供了相当强的灵活性,我们可以定义一些策略,比如当达成某个条件的时候转换写入文件,新写入文件的名字,写入时候的分隔符等等。 如果选择使用的话,Storm 有提供部分接口供我们使用,但如果我们觉得不够丰富也可以自定义相应的类。下面我们看看如何控制这些策略吧。 RecordFormat 这是一个接口,允许你自由定义接收到内容的格式。 public interface RecordFormat extends Serializable { byte[] format(Tuple tuple); } Storm 提供了 DelimitedRecordFormat ,使用方法在上面已经有了。这个类默认的分割符是逗号",",而你可以通过 withFieldDelimiter 方法改变分隔符。如果你的初始分隔符不是逗号的话,那么也可以重写写一个类实现 RecordFormat 接口即可。 FileNameFormat 同样是一个接口。 public interface FileNameFormat extends Serializable { void prepare(Map conf, TopologyContext topologyContext); String getName(long rotation, long timeStamp); String getPath(); } Storm 所提供的默认的是 org.apache.storm.hdfs.format.DefaultFileNameFormat 。默认人使用的转换文件名有点长,格式是这样的: {prefix}{componentId}-{taskId}-{rotationNum}-{timestamp}{extension} 例如: MyBolt-5-7-1390579837830.txt 默认情况下,前缀是空的,扩展标识是".txt"。 SyncPolicy 同步策略允许你将 buffered data 缓冲到 Hdfs 文件中(从而client可以读取数据),通过实现org.apache.storm.hdfs.sync.SyncPolicy 接口: public interface SyncPolicy extends Serializable { boolean mark(Tuple tuple, long offset); void reset(); } FileRotationPolicy 这个接口允许你控制什么情况下转换写入文件。 public interface FileRotationPolicy extends Serializable { boolean mark(Tuple tuple, long offset); void reset(); } Storm 有提供三个实现该接口的类: 最简单的就是不进行转换的org.apache.storm.hdfs.bolt.rotation.NoRotationPolicy ,就是什么也不干。 通过文件大小触发转换的 org.apache.storm.hdfs.bolt.rotation.FileSizeRotationPolicy。 通过时间条件来触发转换的 org.apache.storm.hdfs.bolt.rotation.TimedRotationPolicy。 如果有更加复杂的需求也可以自己定义。 RotationAction 这个主要是提供一个或多个 hook ,可加可不加。主要是在触发写入文件转换的时候会启动。 public interface RotationAction extends Serializable { void execute(FileSystem fileSystem, Path filePath) throws IOException; } 三.实现一个例子 了解了上面的情况后,我们会实现一个例子,根据写入记录的多少来控制写入转换(改变写入的文件),并且转换后文件的名字表示当前是第几次转换。 首先来看看 HdfsBolt 的内容: RecordFormat format = new DelimitedRecordFormat().withFieldDelimiter(" "); // sync the filesystem after every 1k tuples SyncPolicy syncPolicy = new CountSyncPolicy(1000); // FileRotationPolicy rotationPolicy = new FileSizeRotationPolicy(1.0f, FileSizeRotationPolicy.Units.KB); /** rotate file with Date,every month create a new file * format:yyyymm.txt */ FileRotationPolicy rotationPolicy = new CountStrRotationPolicy(); FileNameFormat fileNameFormat = new TimesFileNameFormat().withPath("/test/"); RotationAction action = new NewFileAction(); HdfsBolt bolt = new HdfsBolt() .withFsUrl("hdfs://127.0.0.1:9000") .withFileNameFormat(fileNameFormat) .withRecordFormat(format) .withRotationPolicy(rotationPolicy) .withSyncPolicy(syncPolicy) .addRotationAction(action); 然后分别来看各个策略的类。 FileRotationPolicy import org.apache.storm.hdfs.bolt.rotation.FileRotationPolicy; import org.apache.storm.tuple.Tuple; import java.text.SimpleDateFormat; import java.util.Date; /** * 计数以改变Hdfs写入文件的位置,当写入10次的时候,则更改写入文件,更改名字取决于 “TimesFileNameFormat” * 这个类是线程安全 */ public class CountStrRotationPolicy implements FileRotationPolicy { private SimpleDateFormat df = new SimpleDateFormat("yyyyMM"); private String date = null; private int count = 0; public CountStrRotationPolicy(){ this.date = df.format(new Date()); // this.date = df.format(new Date()); } /** * Called for every tuple the HdfsBolt executes. * * @param tuple The tuple executed. * @param offset current offset of file being written * @return true if a file rotation should be performed */ @Override public boolean mark(Tuple tuple, long offset) { count ++; if(count == 10) { System.out.print("num :" +count + " "); count = 0; return true; } else { return false; } } /** * Called after the HdfsBolt rotates a file. */ @Override public void reset() { } @Override public FileRotationPolicy copy() { return new CountStrRotationPolicy(); } } FileNameFormat import org.apache.storm.hdfs.bolt.format.FileNameFormat; import org.apache.storm.task.TopologyContext; import java.util.Map; /** * 决定重新写入文件时候的名字 * 这里会返回是第几次转换写入文件,将这个第几次做为文件名 */ public class TimesFileNameFormat implements FileNameFormat { //默认路径 private String path = "/storm"; //默认后缀 private String extension = ".txt"; private Long times = new Long(0); public TimesFileNameFormat withPath(String path){ this.path = path; return this; } @Override public void prepare(Map conf, TopologyContext topologyContext) { } @Override public String getName(long rotation, long timeStamp) { times ++ ; //返回文件名,文件名为更换写入文件次数 return times.toString() + this.extension; } public String getPath(){ return this.path; } } RotationAction import org.apache.hadoop.fs.FileContext; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.storm.hdfs.common.rotation.RotationAction; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; import java.net.URI; /** 当转换写入文件时候调用的 hook ,这里仅写入日志。 */ public class NewFileAction implements RotationAction { private static final Logger LOG = LoggerFactory.getLogger(NewFileAction.class); @Override public void execute(FileSystem fileSystem, Path filePath) throws IOException { LOG.info("Hdfs change the written file!!"); return; } } OK,这样就大功告成了。通过上面的代码,每接收到 10 个 Tuple 后就会转换写入文件,新文件的名字就是第几次转换。 完整代码包括一个随机生成字符串的 Spout ,可以到我的 github 上查看。 StormHdfsDemo:https://github.com/shezhiming/StormHdfsDemo 更多干货,欢迎关注公众号,哈尔的数据城堡。

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

苹果新的编程语言 Swift 语言进阶(五)--控制流

Swift 语言支持C语言全部的控制语句。包含for 和while循环语句,if和switch条件语句,以及break和continue控制语句等。 Swift 语言除了支持以上语句,还添加了一个for-in循环语句。来更方面地遍历数组、词典、范围、字符串和其他序列等。 1、for-in循环 forindexin1...5{ println("\(index)times 5 is\(index*5)") } 以上for-in循环用来遍历一个闭合的的范围。 为了语句的简洁。以上语句中使用到的index能够从循环运行体中判断出是一个常量。因此该常量不须要在使用之前使用let keyword来声明。 假设想使用它作为变量。就必须对它进行声明。 在循环语句或条件语句中声明的常量或变量仅在循环运行体或循环运行体中有效。 假设不须要使用for-in范围的每个值,以上语句还能够採用例如以下形式: letbase=3 letpower=10 varanswer=1 for_in1...power{ answer*=base } 2 Switch语句 Swift 语言对switch语法进行了优化,功能做了增强。 优化后的switch 语法更加安全、语义更加清楚。 Swift 语言要求switch的每一个分支必须是全面的。即switch声明的每一个可能的值必须运行Case分支之中的一个。假设每一个Case分支已经全面覆盖了每种情况,default 分支语句就能够省去。 Swift 语言不同意带有空运行体的Case分支, 每一个Case运行体必须相应一个Case分支,但多个可能的值能够放到一个Case声明中,用于相应一个运行分支,一个Case分支的多个匹配值由逗号切割,相应一个Case分支的多个匹配值能够分成多个行显示。 Swift 语言不须要在每一个Case分支加入一个多余的break语句。Swift 语言运行完一个Case分支的运行体后,自己主动退出switch语句,这样能够避免C语言常常出现的因为缺少break语句引起的逻辑错误。 Swift 语言支持使用Fallthrough语句来明白说明在一个Case运行体运行完后不退出switch语句而是直接运行接着的Case运行体或者default运行体。 letsomeCharacter:Character="e" switchsomeCharacter{ case"a","e","i","o","u": println("\(someCharacter)is a vowel") case"b","c","d","f","g","h","j","k","l","m", "n","p","q","r","s","t","v","w","x","y","z": println("\(someCharacter)is a consonant") default: println("\(someCharacter)is not a vowel or a consonant") } 因为Swift 语言要求每一个Case分支必须包括一个至少一条语句的运行体。例如以下代码因为第一个case分支 缺少运行体将报一个编译错误,该优化也从语法上避免了一个case运行另外的case的情况,语法也更清晰。 letanotherCharacter:Character="a" switchanotherCharacter{ case"a": case"A": println("The letter A") default: println("Not the letter A") } Swift 语言 的switchCase 分支能够採用很多不同类型的匹配模式,包含范围、多元组。 例如以下是一个使用多元组匹配的样例。 letsomePoint= (1,1) switchsomePoint{ case(0,0): println("(0, 0) is at the origin") case(_,0): println("(\(somePoint.0), 0) is on the x-axis") case(0,_): println("(0,\(somePoint.1)) is on the y-axis") case(-2...2, -2...2): println("(\(somePoint.0),\(somePoint.1)) is inside the box") default: println("(\(somePoint.0),\(somePoint.1)) is outside of the box") } Swift 语言同意多个case 分支符合某个条件。在某个值匹配多个case 分支的情况下, Swift 语言规定总是使用第一个最先匹配的分支。 如以上样例尽管(0, 0)点匹配全部四个case 分支。但它仅仅运行首先匹配到的分支case (0,0)相应的运行体,其他匹配的分支将被忽略。 下面是一个使用范围进行匹配的样例。 letcount=3_000 varnaturalCount:String switchcount{ case0: naturalCount="no" case1...3: naturalCount="a few" case4...9: naturalCount="several" case10...99: naturalCount="tens of" case100...999: naturalCount="hundreds of" case1000...999_999: naturalCount="thousands of" default: naturalCount="millions and millions of" } 在Case分支中。匹配值还能被绑定到一个暂时常量或变量,以便在也仅仅能在Case的运行体中使用。 例如以下是一个使用值绑定的样例。 letanotherPoint= (2,0) switchanotherPoint{ case(letx,0): println("on the x-axis with an x value of\(x)") case(0,lety): println("on the y-axis with a y value of\(y)") caselet(x,y): println("somewhere else at (\(x),\(y))") } 每个Case分支还能使用where从句来检查额外的更加复杂的条件。例如以下所看到的: letyetAnotherPoint= (1, -1) switchyetAnotherPoint{ caselet(x,y)wherex==y: println("(\(x),\(y)) is on the line x == y") caselet(x,y)wherex== -y: println("(\(x),\(y)) is on the line x == -y") caselet(x,y): println("(\(x),\(y)) is just some arbitrary point") } 3、传输控制语句和标签语句 Swift 除了支持标准的continue、break、return传输控制语句外,还提供了一个以上提到的fallthrough传输控制语句。 Swift 也支持循环语句或switch语句的嵌套,另外Swift 还支持为一个循环语句或switch语句加入一个标签,然后能够使用传输控制语句continue、break跳转到该标签语句处运行。例如以下所看到的: label name:whilecondition{ switchcondition{ casevalue 1: statements 1 breaklabel name casevalue 2: statements 2 continuelabel name } } 版权全部。请转载时清楚注明链接和出处! 本文转自mfrbuaa博客园博客,原文链接:http://www.cnblogs.com/mfrbuaa/p/5132612.html,如需转载请自行联系原作者

资源下载

更多资源
Spring

Spring

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

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

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部分的功能。

用户登录
用户注册