首页 文章 精选 留言 我的

精选列表

搜索[文档识别],共10000篇文章
优秀的个人博客,低调大师

Redisson官方文档 - 13. 工具

13.1. 集群管理工具 Redisson集群管理工具提供了通过程序化的方式,像redis-trib.rb脚本一样方便地管理Redis集群的工具。 13.1.1 创建集群 以下范例展示了如何创建三主三从的Redis集群。 ClusterNodes clusterNodes = ClusterNodes.create() .master("127.0.0.1:7000").withSlaves("127.0.0.1:7001", "127.0.0.1:7002") .master("127.0.0.1:7003").withSlaves("127.0.0.1:7004") .master("127.0.0.1:7005"); ClusterManagementTool.createCluster(clusterNodes); 主节点127.0.0.1:7000的从节点有127.0.0.1:7001和127.0.0.1:7002。 主节点127.0.0.1:7003的从节点是127.0.0.1:7003。 主节点127.0.0.1:7005没有从节点。 13.1.2 踢出节点 以下范例展示了如何将一个节点踢出集群。 ClusterManagementTool.removeNode("127.0.0.1:7000", "127.0.0.1:7002"); // 或 redisson.getClusterNodesGroup().removeNode("127.0.0.1:7002"); 将从节点127.0.0.1:7002从其主节点127.0.0.1:7000里踢出。 13.1.3 数据槽迁移 以下范例展示了如何将数据槽在集群的主节点之间迁移。 ClusterManagementTool.moveSlots("127.0.0.1:7000", "127.0.0.1:7002", 23, 419, 4712, 8490); // 或 redisson.getClusterNodesGroup().moveSlots("127.0.0.1:7000", "127.0.0.1:7002", 23, 419, 4712, 8490); 将番号为23,419,4712和8490的数据槽从127.0.0.1:7002节点迁移至127.0.0.1:7000节点。 以下范例展示了如何将一个范围的数据槽在集群的主节点之间迁移。 ClusterManagementTool.moveSlotsRange("127.0.0.1:7000", "127.0.0.1:7002", 51, 9811); // 或 redisson.getClusterNodesGroup().moveSlotsRange("127.0.0.1:7000", "127.0.0.1:7002", 51, 9811); 将番号范围在[51, 9811](含)之间的数据槽从127.0.0.1:7002节点移动到127.0.0.1:7000节点。 13.1.4 添加从节点 以下范例展示了如何向集群中添加从节点。 ClusterManagementTool.addSlaveNode("127.0.0.1:7000", "127.0.0.1:7003"); // 或 redisson.getClusterNodesGroup().addSlaveNode("127.0.0.1:7003"); 将127.0.0.1:7003作为从节点添加至127.0.0.1:7000所在的集群里。 13.2.5 添加主节点 以下范例展示了如何向集群中添加主节点。 ClusterManagementTool.addMasterNode("127.0.0.1:7000", "127.0.0.1:7004"); // 或 redisson.getClusterNodesGroup().addMasterNode("127.0.0.1:7004"); 将127.0.0.1:7004作为主节点添加至127.0.0.1:7000所在的集群里。Adds master node 127.0.0.1:7004 to cluster where 127.0.0.1:7000 participate in 该功能仅限于Redisson PRO版本。

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

Redisson官方文档 - 1. 概述

Redisson是一个在Redis的基础上实现的Java驻内存数据网格(In-Memory Data Grid)。它不仅提供了一系列的分布式的Java常用对象,还提供了许多分布式服务。其中包括(BitSet, Set, Multimap, SortedSet, Map, List, Queue, BlockingQueue, Deque, BlockingDeque, Semaphore, Lock, AtomicLong, CountDownLatch, Publish / Subscribe, Bloom filter, Remote service, Spring cache, Executor service, Live Object service, Scheduler service) Redisson提供了使用Redis的最简单和最便捷的方法。Redisson的宗旨是促进使用者对Redis的关注分离(Separation of Concern),从而让使用者能够将精力更集中地放在处理业务逻辑上。 关于Redisson项目的详细介绍可以在官方网站找到。 每个Redis服务实例都能管理多达1TB的内存。 能够完美的在云计算环境里使用,并且支持AWS ElastiCache主备版,AWS ElastiCache集群版,Azure Redis Cache和阿里云(Aliyun)的云数据库Redis版 以下是Redisson的结构: Redisson作为独立节点 可以用于独立执行其他节点发布到分布式执行服务 和 分布式调度任务服务 里的远程任务。 如果你现在正在使用其他的Redis的Java客户端,那么Redis命令和Redisson对象匹配列表 能够帮助你轻松的将现有代码迁徙到Redisson框架里来。 Redisson底层采用的是Netty 框架。支持Redis 2.8以上版本,支持Java1.6+以上版本。

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

Apache Storm 官方文档 —— 常用模式

本文列出了 Storm 拓扑中使用的一些常见模式,包括: 数据流的 join 批处理 BasicBolt 内存缓存与域分组的结合 Top N 流式计算 TimeCacheMap CoordinatedBolt 与 KeyedFairBolt Joins 数据流的 join 一般指的是通过共有的域来聚合两个或多个数据流的过程。与一般的数据库中 join 操作要求有限的输入与清晰的语义不同,数据流 join 的输入往往是无限的数据集,而且并不具备明确的语义。 join 的类型一般是由应用的需求决定的。有些应用需要将两个流在某个固定时间内的所有 tuple 进行 join,另外一些应用却可能要求对每个 join 域的 join 操作过程的两侧只保留一个 tuple,而其他的应用也许还有一些其他需求。不过这些 join 类型一般都会有一个基本的模式,那就是将多个输入流进行分区。Storm 可以很容易地使用域分组的方法将多个输入流聚集到一个联结 bolt 中,比如下面这样: builder.setBolt("join", new MyJoiner(), parallelism) .fieldsGrouping("1", new Fields("joinfield1", "joinfield2")) .fieldsGrouping("2", new Fields("joinfield1", "joinfield2")) .fieldsGrouping("3", new Fields("joinfield1", "joinfield2")); 当然,上面的代码只是个例子,实际上不同的流完全可以具有不同的输入域。 批处理 通常由于效率或者其他方面的原因,你需要使用将 tuple 们组合成 batch 来处理,而不是一个个分别处理它们。比如,在做数据库更新操作或者流聚合操作时,你就会需要这样的批处理形式。 要确保数据处理的可靠性,正确的方式是在 bolt 进行批处理之前将 tuple 们缓存在一个实例变量中。在完成批处理操作之后,你就可以一起 ack 所有的缓存的 tuple 了。 如果这个批处理 bolt 还需要继续向下游发送 tuple,你可能还需要使用多锚定(multi-anchoring)来确保可靠性。具体怎么做取决于应用的需求。想要了解更多关于可靠性的工作机制的内容请参考消息的可靠性保障一文。 BasicBolt Bolt 处理 tuple 的一种基本模式是在execute方法中读取输入 tuple、发送出基于输入 tuple 的新 tuple,然后在方法末尾对 tuple 进行应答(ack)。符合这种模式的 bolt 一般是一种函数或者过滤器。对于这种基本的处理模式,Storm 提供了IBasicBolt接口来自动实现这个过程。更多内容请参考消息的可靠性保障一文。 内存缓存与域分组的结合 在 Storm 的 bolt 中保存一定的缓存也是一种比较常见的方式。尤其是在于域分组结合的时候,缓存的作用特别显著。例如,假如你有一个用于将短链接(short URLs,例如 bit.ly, t.co,等等)转化成长链接(龙 URLs)的 bolt。你可以通过一个将短链接映射到长链接的 LRU 缓存来提高系统的性能,避免反复的 HTTP 请求操作。假如现在有一个名为 “urls” 的组件用于发送短链接,另外有一个 “expand” 组件用于将短链接扩展为长链接,并且在 “expand” 内部保留一个缓存。让我们来看看下面两段代码有什么不同: builder.setBolt("expand", new ExpandUrl(), parallelism) .shuffleGrouping(1); builder.setBolt("expand", new ExpandUrl(), parallelism) .fieldsGrouping("urls", new Fields("url")); 由于域分组可以使得相同的 URL 永远被发往同一个 task,第二段代码会比第一段代码高效得多。这样可以避免在不同的 task 的缓存中的复制动作,并且看上去短 URL 可以更好地在命中缓存。 Top N Storm 中一种常见的连续计算模式是计算数据流中某种形式的 Top N 结果。假如现在有一个可以以 [“value”, “count”] 的形式发送 tuple 的 bolt,并且你需要一个可以根据 count 计算结果输出前 N 个 tuple 的 bolt。实现这个操作的最简单的方法就是使用一个对数据流进行全局分组的 bolt,并且在内存中维护一个包含 top N 结果的列表。 这种方法并不适用于大规模数据流,因为整个数据流都会发往同一个 task,会造成该 task 的内存负载过高。更好的做法是将数据流分区,同时对每个分区计算 top N 结果,然后将这些结果汇总来得到最终的全局 top N 结果。下面是这个模式的代码: builder.setBolt("rank", new RankObjects(), parallelism) .fieldsGrouping("objects", new Fields("value")); builder.setBolt("merge", new MergeObjects()) .globalGrouping("rank"); 这个方法之所以可行是因为第一个 bolt 的域分组操作确保了每个小分区在语义上的正确性。你可以在storm-starter里看到使用这个模式的一个例子。 当然,如果待处理的数据集存在较严重的数据倾斜,那么还是应该使用 partialKeyGrouping 来代替 fieldsGrouping,因为 partialKeyGrouping 可以通过两个下游 bolt 分散每个 key 的负载。 builder.setBolt("count", new CountObjects(), parallelism) .partialKeyGrouping("objects", new Fields("value")); builder.setBolt("rank" new AggregateCountsAndRank(), parallelism) .fieldsGrouping("count", new Fields("key")) builder.setBolt("merge", new MergeRanksObjects()) .globalGrouping("rank"); 这个拓扑中需要一个中间层来聚合来自上游 bolt 数据流的分区计数结果,但这一层仅仅会做一个简单的聚合处理,这样 bolt 就不会受到由于数据倾斜带来的负载压力。你可以在storm-starter中看到使用这个模式的一个例子。 支持 LRU 的 TimeCacheMap 有时候你可能会需要一个能够保留“活跃的”数据并且能够使得超时的“非活跃的”数据自动失效的缓存。TimeCacheMap是一个可以高效地实现此功能的数据结构。它还提供了一个钩子用于实现在数据失效后的回调操作。 用于分布式 RPC 的 CoordinatedBolt 与 KeyedFairBolt 在构建 Storm 上层的分布式 RPC 应用时,通常会用到两种常用的模式。现在这两种模式已经被封装为CoordinatedBolt和KeyedFairBolt,并且已经加入了 Storm 标准库中。 CoordinatedBolt将你的处理逻辑 bolt 包装起来,并且在你的 bolt 收到了指定请求的所有 tuple 之后发出通知。CoordinatedBolt中大量使用了直接数据流组来实现此功能。 KeyedFairBolt同样包装了你的处理逻辑 bolt,并且可以让你的拓扑同时处理多个 DRPC 调用,而不是每次只执行一个。 如果需要了解更多内容请参考分布式 RPC一文。 转载自并发编程网 - ifeve.com

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

Apache Storm 官方文档 —— Trident 教程

Trident 是 Storm 的一种高度抽象的实时计算模型,它可以将高吞吐量(每秒百万级)数据输入、有状态的流式处理与低延时的分布式查询无缝结合起来。如果你了解 Pig 或者 Cascading 这样的高级批处理工具,你就会发现他们和 Trident 的概念非常相似。Trident 同样有联结(join)、聚合(aggregation)、分组(grouping)、函数(function)以及过滤器(filter)这些功能。Trident 为数据库或者其他持久化存储上层的状态化、增量式处理提供了基础原语。由于 Trident 有着一致的、恰好一次的语义,因此推断出 Trident 拓扑的状态也是一件很容易的事。 使用范例 让我们先从一个使用 Trident 的例子开始。这个例子中做了两件事情: 从一个句子的输入数据流中计算出单词流的数量 实现对一个单词列表中每个单词总数的查询 为了实现这个目的,这个例子将会从下面的数据源中无限循环地读取语句数据流: FixedBatchSpout spout = new FixedBatchSpout(new Fields("sentence"), 3, new Values("the cow jumped over the moon"), new Values("the man went to the store and bought some candy"), new Values("four score and seven years ago"), new Values("how many apples can you eat")); spout.setCycle(true); 这个 Spout 会循环地访问语句集来生成语句数据流。下面的代码就是用来实现计算过程中的单词数据流统计部分: TridentTopology topology = new TridentTopology(); TridentState wordCounts = topology.newStream("spout1", spout) .each(new Fields("sentence"), new Split(), new Fields("word")) .groupBy(new Fields("word")) .persistentAggregate(new MemoryMapState.Factory(), new Count(), new Fields("count")) .parallelismHint(6); 让我们一行行地来分析上面的代码。首先我们创建了一个TridentTopology对象,这个对象提供了构造 Trident 计算过程的接口。TridentTopology有一个叫做newStream的方法,这个方法可以从一个输入数据源中读取数据创建一个新的数据流。在这个例子中,输入的数据源就是前面定义的FixedBatchSpout。输入数据源也可以是像 Kestrel 和 Kafka 这样的消息系统。Trident 会通过 ZooKeeper 一直跟踪每个输入数据源的一小部分状态(Trident 具体消费对象的相关元数据)。例如这里的 “spout1” 就对应着 ZooKeeper 中的一个节点,而 Trident 就会在该节点中存放数据源的元数据(metadata)。 Trident 会将数据流处理为很多个小块 tuple 的集合,例如,输入的句子流就会像下面这样被分割成很多个小块: 这些小块的大小主要取决于你的输入吞吐量,一般可能会在数万甚至数百万元组的级别。 Trident 为这些小块提供了一个完全成熟的批处理 API。这个 API 和你见到过的 Pig 或者 Cascading 这样的 Hadoop 的高级抽象语言很相似:你可以处理分组(group by)、联结(join)、聚合(aggregation)、函数(function)、过滤器(filter)等各种操作。当然,分别处理每个小块并不是件好事,所以,Trident 提供了适用于处理各个小块之间的聚合操作的函数,并且可以在聚合后将结果保存到持久化存储中,而且无论是内存、Memcached、Cassandra 还是其他类型的存储都可以支持。最后,Trident 还提供了用于查询实时状态结果的一级接口。而这个结果状态既可以像这个例子中演示的那样由 Trident 负责更新,也可以作为一个独立的状态数据源而存在。 再回到这个例子中,输入数据源 spout 发送出了一个名为 “sentence” 的数据流。接下来拓扑中定义了一个Split方法用于处理流中的每个 tuple,这个方法接收 “sentence” 域并将其分割成若干个单词。每个 sentence tuple 都会创建很多个单词 tuple —— 例如 “the cow jumped over the moon” 这个句子就会创建 6 个 “word” tuple,下面是Split的定义: public class Split extends BaseFunction { public void execute(TridentTuple tuple, TridentCollector collector) { String sentence = tuple.getString(0); for(String word: sentence.split(" ")) { collector.emit(new Values(word)); } } } 从上面的代码中你会发现这个过程真的很简单。这个方法中的所有操作仅仅是抓取句子、以空格分隔句子并且为每个单词发射一个 tuple。 拓扑的剩余部分负责统计单词的数量并将结果保存到持久化存储中。首先,数据流根据 “word” 域分组,然后使用Count聚合器持续聚合每个小组。persistentAggregate方法用于存储并更新 state 源中的聚合结果。在这个例子中,单词的数量结果是保存在内存中的,不过可以根据需要切换到 Memcached、Cassandra 或者其他持久化存储中。切换存储模型也非常简单,只需要像下面这样(使用trident-memcached修改persistentAggregate行中的一个参数(其中,“serverLocations” 是 Memcached 集群的地址/端口列表)即可: .persistentAggregate(MemcachedState.transactional(serverLocations), new Count(), new Fields("count")) persistentAggregate方法所存储的值就表示所有从数据流中发送出来的块的聚合结果。 Trident 的另一个很酷的特性就是它支持完全容错性和恰好一次处理的语义。如果处理过程中出现错误需要重新执行处理操作,Trident 不会向数据库中提交多次来自相同的源数据的更新操作,这就是 Trident 持久化 state 的方式。 persistentAggregate方法也可以将数据流结果传入一个TridentState对象中。这种情况下,这个TridentState就表示所有的单词统计信息。这样我们就可以使用TridentState对象来实现整个计算过程中的分布式查询部分。 接下来我们就可以在拓扑中实现 word count 的一个低延时分布式查询。这个查询接收一个由空格分隔的单词列表作为参数,然后返回这些单词的数量统计结果。这个查询看上去与普通的 RPC 调用并没有什么分别,不过在后台他们是并发执行的。下面是一个实现这种查询的例子: DRPCClient client = new DRPCClient("drpc.server.location", 3772); System.out.println(client.execute("words", "cat dog the man"); // prints the JSON-encoded result, e.g.: "[[5078]]" 如你所见,这个查询看上去只是一个普通的远程过程调用(RPC),不过在后台他是在一个 Storm 集群中并发执行的。这种查询的端到端延时一般在 10 ms 左右。当然,更大量的查询会花费更长的时间,尽管这些查询还是取决于你为这个计算过程分配了多少时间。 拓扑中的分布式查询的实现是这样的: topology.newDRPCStream("words") .each(new Fields("args"), new Split(), new Fields("word")) .groupBy(new Fields("word")) .stateQuery(wordCounts, new Fields("word"), new MapGet(), new Fields("count")) .each(new Fields("count"), new FilterNull()) .aggregate(new Fields("count"), new Sum(), new Fields("sum")); 这里还需要使用前面的TridentTopology对象来创建一个 DRPC 数据流,这个创建数据流的方法叫做 “words”。前面使用DRPCClient进行 RPC 调用的第一个参数必须与这个方法名完全相同。 在这段代码里,首先是使用Split方法来将请求的参数分割成若干个单词。这些单词构成的单词流是通过 “word” 域来分组的,而stateQuery运算符就是用来查询拓扑中第一个部分中生成的TridentState对象的。stateQuery接收一个 state(在这个例子中就是拓扑前面计算得到的单词数结果)和查询这个 state 的方法作为参数。在这个例子里,stateQuery调用了MapGet方法,用于获取每个单词的个数。由于 DRPC 数据流是和 TridentState 采用的完全相同的方式进行分组的(通过 “word” 域),每个单词查询都可以精确地定位到 TridentState 对象中的指定部分,同时 TridentState 对象中维护着对应的单词的更新状态。 接下来,个数为 0的单词会被FilterNull过滤器过滤掉,然后就可以使用Sum聚合器来获取其他的单词统计个数。接着 Trident 就会自动将结果返回给等待的客户端。 Trident 很聪明,它知道怎么以最好的性能运行拓扑。在这个拓扑中还有两个会自动发生的有趣的事: 从 state 中读取或写入的操作(例如 persistentAggregate 和 stateQuery)会自动批处理化。因此,如果当前的批处理过程需要对数据库执行 20 个更新操作,Trident 就会自动将读取或写入操作当作批处理过程,仅仅会对数据库发送一次读请求和一次写请求,而不是发送 20 次读请求和 20 次写请求(而且一般你还可以在你的 state 里使用缓存来消除读请求)。这样做就有两个方面的好处:可以按照你指定的方式来执行你的计算过程,同时还可以维持较好的性能。 Trident 的聚合器是高度优化的。在向网络中发送 tuple 之前,Trident 有时候会做部分聚合操作,而不是将一个分组的所有的 tuple 一股脑地发送到同一台机器中来执行聚合。例如,Count聚合器就是这样先计算每个小块的个数,然后向网络中发送很多个部分计数的结果,接着再将所有的部分计数结果汇总来得到最终的统计结果。这个技术与 MapReduce 的 combiner 模型很相似。 我们再来看看 Trident 的另一个例子。 Reach 这个例子是一个纯粹的 DRPC 拓扑,计算了一个指定 URL 的 Reach 数。Reach 指的是 Twitter 上能够看到一个指定的 URL 的独立用户数。要想计算 Reach,你需要先提取所有转发了该 URL 的用户,提取这些用户的关注者,将关注者放入一个 set 集合中来去除重复的关注者,然后再统计这个 set 中的数量。对于单一的一台机器来说,计算 reach 太耗时了,这个过程大概需要数千次数据库调用并生成数千万 tuple。而使用 Storm 和 Trident 就可以通过一个集群来将计算过程的每个步骤进行并行化处理。 这个拓扑会从两个 state 源中读取数据。其中一个数据库建立了 URL 和转发了该 URL 的用户列表的关联表。另一个数据库中建立了用户和用户的关注者列表的关联表。拓扑的定义是这样的: TridentState urlToTweeters = topology.newStaticState(getUrlToTweetersState()); TridentState tweetersToFollowers = topology.newStaticState(getTweeterToFollowersState()); topology.newDRPCStream("reach") .stateQuery(urlToTweeters, new Fields("args"), new MapGet(), new Fields("tweeters")) .each(new Fields("tweeters"), new ExpandList(), new Fields("tweeter")) .shuffle() .stateQuery(tweetersToFollowers, new Fields("tweeter"), new MapGet(), new Fields("followers")) .parallelismHint(200) .each(new Fields("followers"), new ExpandList(), new Fields("follower")) .groupBy(new Fields("follower")) .aggregate(new One(), new Fields("one")) .parallelismHint(20) .aggregate(new Count(), new Fields("reach")); 这个拓扑使用newStaticState方法创建了两个分别对应外部于两个外部数据库的TridentState对象。在拓扑的后续部分就可以对这两个TridentState对象执行查询操作。和 state 的所有数据源一样,为了最大程度地提升效率,对这些数据库的查询将会自动地批处理化。 拓扑的定义很直接 —— 就是一个简单的批处理 job。首先,会通过查询 urlToTweeters 数据库来获取转发了 URL 的用户列表,然后就可以调用ExpandList方法来为每个 tweeter 创建一个 tuple。 接下来必须要获取每个 tweeter 的关注者。由于需要调用 shuffle 方法将所有的 tweeter 均衡分配到拓扑的所有 worker 中,所以这个步骤必须并发进行,这一点非常重要。然后就可以查询关注者数据库来获取每个 tweeter 的关注者列表。你可能注意到了这个过程的并行度非常高,因为这是整个计算过程中复杂度最高的部分。 再接下来,关注者就会被放入一个单独的 set 集合中用于计数。这里包含两个步骤。首先,会根据 “follower” 域来执行 “group by” 分组操作,并在每个组上运行One聚合器。“One”聚合器的作用仅仅是为每个组发送一个包含数字 1 的 tuple。然后,就可以通过统计这些 one 结果来得到关注者 set 的大小,也就是真正的关注者数量。下面是 “One” 聚合器的定义: public class One implements CombinerAggregator<Integer> { public Integer init(TridentTuple tuple) { return 1; } public Integer combine(Integer val1, Integer val2) { return 1; } public Integer zero() { return 1; } } 这是一个“组合聚合器”,它知道怎样在向网络中发送 tuple 之前以最好的效率进行部分聚合操作。同样,Sum 也是一个组合聚合器,所以在拓扑结尾的全局统计操作也会有很高的效率。 下面让我们再来看看 Trident 中的一些细节。 域(Fields)与元组(tuples) Trident 的数据模型 TridentTuple 是一个指定的值列表。在一个拓扑中,tuple 是在一系列操作中不断生成的。这些操作一般会输入一个“输入域”(input fields)集合,然后发送出一个“方法域”(function fields)的集合。输入域主要用于选取一个 tuple 的子集作为操作的输入,而“方法域”主要用于为该操作的输出结果域命名。 我们来看看这样一个场景。假设你有一个名为 “stream” 的数据流,其中包含域 “x”、“y” 和 “z”。如果要运行一个接收 “y” 作为输入的过滤器 MyFilter,你可以这样写: stream.each(new Fields("y"), new MyFilter()) 再假设 MyFilter 的实现是这样的: public class MyFilter extends BaseFilter { public boolean isKeep(TridentTuple tuple) { return tuple.getInteger(0) < 10; } } 这样就会保留所有 “y” 域的值小于 10 的 tuple。MyFilter 输入的 TridentTuple 将会仅包含有 “y” 域。值得注意的是,Trident 可以在选取输入域时以一种非常高效的方式来投射 tuple 的子集:这个投射过程非常灵活。 我们再来看看 “function fields” 是怎么工作的。假设你有这样一个函数: public class AddAndMultiply extends BaseFunction { public void execute(TridentTuple tuple, TridentCollector collector) { int i1 = tuple.getInteger(0); int i2 = tuple.getInteger(1); collector.emit(new Values(i1 + i2, i1 * i2)); } } 这个函数接收两个数字作为输入,然后发送出两个新值:分别是两个数字的和和乘积。再假定你有一个包含 “x”、“y” 和 “z” 域的数据流,你可以这样使用这个函数: stream.each(new Fields("x", "y"), new AddAndMultiply(), new Fields("added", "multiplied")); 这个函数的输出增加了两个新的域。因此,这个 each 调用的输出 tuple 会包含 5 个域:“x”、“y” 、“z”、“added” 和 “multiplied”。其中 “added” 与 AddAndMultiply 的第一个输出值相对应,“multiplied” 和 AddAndMultiply 的第二个输出值相对应。 另一方面,通过聚合器,函数域也可以替换输入 tuple 的域。假如你有一个包含域 “val1” 和域 “val2” 的数据流,通过这样的操作: stream.aggregate(new Fields("val2"), new Sum(), new Fields("sum")) 就会使得输出数据流中只包含一个只带有 “sum” 的域的 tuple,这个 “sum” 域就代表了在哪个批处理块中所有的 “val2” 域的总和值。 通过数据流分组,输出就可以同时包含用于分组的域以及由聚合器发送的域。举个例子: stream.groupBy(new Fields("val1")) .aggregate(new Fields("val2"), new Sum(), new Fields("sum")) 这个操作就会使得输出同时包含域 “val1” 以及域 “sum”。 State 实时计算的一个关键问题就在于如何管理状态(state),使得在失败与重试操作之后的更新过程仍然是幂等的。错误是不可消除的,所以在出现节点故障或者其他问题发生时批处理操作还需要进行重试。不过这里最大的问题就在于怎样执行一种合适的状态更新操作(不管是针对外部数据库还是拓扑内部的状态),来使得每个消息都能够被执行且仅仅被执行一次。 这个问题很麻烦,接下来的例子里面就有这样的问题。假如你正在对你的数据流做一个计数聚合操作,并且打算将计数结果存储到一个数据库中。如果你仅仅把计数结果存到数据库里就完事了的话,那么在你继续准备更新某个块的状态的时候,你没法知道到底这个状态有没有被更新过。这个数据块有可能在更新数据库的步骤上成功了,但在后续的步骤中失败了,也有可能先失败了,没有进行更新数据库的操作。你完全不知道到底发生了什么。 Trident 通过下面两件事情解决了这个问题: 在 Trident 中为每个数据块标记了一个唯一的 id,这个 id 就叫做“事务 id”(transaction id)。如果数据块由于失败回滚了,那么它持有的事务 id 不会改变。 State 的更新操作是按照数据块的顺序进行的。也就是说,在成功执行完块 2 的更新操作之前,不会执行块 3 的更新操作。 基于这两个基础特性,你的 state 更新就可以实现恰好一次(exactly-once)的语义。与仅仅向数据库中存储计数不同,这里你可以以一个原子操作的形式把事务 id 和计数值一起存入数据库。在后续更新这个计数值的时候你就可以先比对这个数据块的事务 id。如果比对结果是相同的,那么就可以跳过更新操作 —— 由于 state 的强有序性,可以确定数据库中已经包含有当前数据库的额值。而如果比对结果不同,就可以放心地更新计数值了。 当然,你不需要在拓扑中手动进行这个操作,操作逻辑已经在 State 中封装好了,这个过程会自动进行。同样的,你的 State 对象也不一定要实现事务 id 标记:如果你不想在数据库里耗费空间存储事务 id,你就不用那么做。在这样的情况下,State 会在出现失败的情形下保持“至少处理一次”的操作语义(这样对你的应用也是一件好事)。在这篇文章里你可以了解到更多关于如何实现 State 以及各种容错性权衡技术。 你可以使用任何一种你想要的方法来实现 state 的存储操作。你可以把 state 存入外部数据库,也可以保存在内存中然后在存入 HDFS 中(有点像 HBase 的工作机制)。State 也并不需要一直保存某个状态值。比如,你可以实现一个只保存过去几个小时数据并将其余的数据删除的 State。这是一个实现 State 的例子:Memcached integration。 Trident 拓扑的运行 Trident 拓扑会被编译成一种尽可能和普通拓扑有着同样的运行效率的形式。只有在请求数据的重新分配(比如 groupBy 或者 shuffle 操作)时 tuple 才会被发送到网络中。因此,像下面这样的 Trident 拓扑: 就会被编译成若干个 spout/bolt: 总结 Trident 让实时计算变得非常简单。从上面的描述中,你已经看到了高吞吐量的数据流处理、状态操作以及低延时查询处理是怎样通过 Trident 的 API 来实现无缝结合的。总而言之,Trident 可以让你以一种更加自然,同时仍然保持着良好的性能的方式来实现实时计算。 转载自并发编程网 - ifeve.com

资源下载

更多资源
腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Nacos

Nacos

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

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等操作系统。

用户登录
用户注册