首页 文章 精选 留言 我的

精选列表

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

Apache Storm 官方文档 —— 问题与解决

本文介绍了用户在使用 Storm 过程中遇到的问题与相应的解决方法。 Worker 进程在启动时挂掉而没有留下堆栈跟踪信息的问题 可能出现的现象: 拓扑在一个节点上运行正常,但是多个 worker 进程在多个节点上就会崩溃 解决方案: 你的网络配置可能有问题,导致每个节点无法根据 hostname 连接到其他的节点。ZeroMQ 有时会在不能识别 host 的时候挂掉 进程。如果是这种情况,有两种可行的解决方案: 在 /etc/hosts 文件中配置好 hostname 与 IP 的对应关系 设置一个局域网 DNS 服务器,使得节点可以根据 hostname 定位到其他节点 节点之间无法通信 可能出现的现象: 每个 spout tuple 的处理都不成功 拓扑中的处理过程不起作用 解决方案: Storm 不支持 ipv6,你可以在 supervisor 的 child-opts 配置中添加-Djava.net.preferIPv4Stack=true参数,然后重启 supervisor。 你的网络配置可能存在问题,请参考上个问题中的解决方案。 拓扑在一段时间后停止了 tuple 的处理过程 可能出现的现象: 拓扑正常运行一段时间后突然停止了数据处理过程,并且 spout 的 tuple 一起开始处理失败 解决方案: 这是 ZeroMQ 2.1.10 中的一个已经确认的问题,请将 ZMQ 降级到 2.1.7 版本。 Storm UI 中没有显示出所有的 supervisor 信息 可能出现的现象: Storm UI 中缺少部分 supervisor 的信息 在刷新 Storm UI 页面后 supervisor 列表会变化 解决方案: 确保 supervisor 的本地工作目录是相互独立的(也就是说不要出现在 NFS 中共享同一个目录的情况) 尝试删除 supervisor 的本地工作目录,然后重启 supervisor 后台进程。supervisor 启动时会为自己创建一个唯一的 id 并存储在本地目录中。如果这个 id 被复制到其他节点中,就会让 Storm 无法确定哪个 supervisor 正在运行(这种情况并不少见,如果需要扩展集群,就很容易出现直接将某个节点的 Storm 文件直接复制到新节点的情况 —— 译者注)。 “Multiple defaults.yaml found” 错误 可能出现的现象: 在使用storm jar命令部署拓扑时出现此错误 解决方案: 你很可能在拓扑的 jar 包中包含了 Storm 自身的 jar 包。注意,在打包拓扑时,请不要将 Storm 自身的 jar 包加入,因为 Storm 已经在它的 classpath 中提供了这些 jar 包。 运行 storm jar 命令时出现 “NoSuchMethorError” 可能出现的现象: 运行storm jar命令时出现奇怪的 “NoSuchMethodError” 解决方案: 这可能是由于你部署拓扑的 Storm 版本与你构建拓扑时使用的 Storm 版本不同。请确保你编译拓扑时使用的 Storm 版本与你运行拓扑的 Storm 客户端版本相同。 Kryo ConcurrentModificationException 可能出现的现象: 系统运行时出现如下的异常堆栈跟踪信息 java.lang.RuntimeException: java.util.ConcurrentModificationException at backtype.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:84) at backtype.storm.utils.DisruptorQueue.consumeBatchWhenAvailable(DisruptorQueue.java:55) at backtype.storm.disruptor$consume_batch_when_available.invoke(disruptor.clj:56) at backtype.storm.disruptor$consume_loop_STAR_$fn__1597.invoke(disruptor.clj:67) at backtype.storm.util$async_loop$fn__465.invoke(util.clj:377) at clojure.lang.AFn.run(AFn.java:24) at java.lang.Thread.run(Thread.java:679) Caused by: java.util.ConcurrentModificationException at java.util.LinkedHashMap$LinkedHashIterator.nextEntry(LinkedHashMap.java:390) at java.util.LinkedHashMap$EntryIterator.next(LinkedHashMap.java:409) at java.util.LinkedHashMap$EntryIterator.next(LinkedHashMap.java:408) at java.util.HashMap.writeObject(HashMap.java:1016) at sun.reflect.GeneratedMethodAccessor17.invoke(Unknown Source) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:616) at java.io.ObjectStreamClass.invokeWriteObject(ObjectStreamClass.java:959) at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1480) at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1416) at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1174) at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:346) at backtype.storm.serialization.SerializableSerializer.write(SerializableSerializer.java:21) at com.esotericsoftware.kryo.Kryo.writeClassAndObject(Kryo.java:554) at com.esotericsoftware.kryo.serializers.CollectionSerializer.write(CollectionSerializer.java:77) at com.esotericsoftware.kryo.serializers.CollectionSerializer.write(CollectionSerializer.java:18) at com.esotericsoftware.kryo.Kryo.writeObject(Kryo.java:472) at backtype.storm.serialization.KryoValuesSerializer.serializeInto(KryoValuesSerializer.java:27) 解决方案: 这个信息表示你在将一个可变的对象作为 tuple 发送出去。你发送到 outputcollector 中的所有对象必须是非可变的。这个错误表明对象在被序列化并发送到网络中时你的 bolt 正在修改这个对象。 Storm 中的 NullPointerException 可能出现的现象: Storm 运行中出现了如下的 NullPointerException java.lang.RuntimeException: java.lang.NullPointerException at backtype.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:84) at backtype.storm.utils.DisruptorQueue.consumeBatchWhenAvailable(DisruptorQueue.java:55) at backtype.storm.disruptor$consume_batch_when_available.invoke(disruptor.clj:56) at backtype.storm.disruptor$consume_loop_STAR_$fn__1596.invoke(disruptor.clj:67) at backtype.storm.util$async_loop$fn__465.invoke(util.clj:377) at clojure.lang.AFn.run(AFn.java:24) at java.lang.Thread.run(Thread.java:662) Caused by: java.lang.NullPointerException at backtype.storm.serialization.KryoTupleSerializer.serialize(KryoTupleSerializer.java:24) at backtype.storm.daemon.worker$mk_transfer_fn$fn__4126$fn__4130.invoke(worker.clj:99) at backtype.storm.util$fast_list_map.invoke(util.clj:771) at backtype.storm.daemon.worker$mk_transfer_fn$fn__4126.invoke(worker.clj:99) at backtype.storm.daemon.executor$start_batch_transfer__GT_worker_handler_BANG_$fn__3904.invoke(executor.clj:205) at backtype.storm.disruptor$clojure_handler$reify__1584.onEvent(disruptor.clj:43) at backtype.storm.utils.DisruptorQueue.consumeBatchToCursor(DisruptorQueue.java:81) ... 6 more 解决方案: 这个问题是由于多个线程同时调用OutputCollector中的方法造成的。Storm 中所有的 emit、ack、fail 方法必须在同一个线程中运行。出现这个问题的一种场景是在一个IBasicBolt中创建了一个独立的线程。由于IBasicBolt会在execute方法调用之后自动调用ack,所以这就会出现多个线程同时使用OutputCollector的情况,进而抛出这个异常。也就是说,在使用IBasicBolt时,所有的消息发送操作必须在同一个线程的execute方法中执行。 转载自并发编程网 - ifeve.com

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

Apache Storm 官方文档 —— Storm 与 Kestrel

本文说明了如何使用 Storm 从 Kestrel 集群中消费数据。 前言 Storm 本教程中使用了storm-kestrel项目和storm-starter项目中的例子。建议读者将这几个项目 clone 到本地,并动手运行其中的例子。 Kestrel 本文假定读者可以如此项目所述在本地运行一个 Kestrel 集群。 Kestrel 服务器与队列 Kestrel 服务中包含有一组消息队列。Kestrel 队列是一种非常简单的消息队列,可以运行于 JVM 上,并使用 memcache 协议(以及一些扩展)与客户端交互。详情可以参考storm-kestrel项目中的KestrelThriftClient类的实现。 每个队列均严格遵循先入先出的规则。为了提高服务性能,数据都是缓存在系统内存中的;不过,只有开头的 128MB 是保存在内存中的。在服务停止的时候,队列的状态会保存到一个日志文件中。 请参阅此文了解更多详细信息。 Kestrel 具有 * 快速 * 小巧 * 持久 * 可靠 等特点。 例如,Twitter 就使用 Kestrel 作为消息系统的核心环节,此文中介绍了相关信息。 ** 向 Kestrel 中添加数据 首先,我们需要一个可以向 Kestrel 的队列添加数据的程序。下述方法使用了storm-kestrel项目中的KestrelClient的实现。该方法从一个包含 5 个句子的数组中随机选择一个句子添加到 Kestrel 的队列中。 private static void queueSentenceItems(KestrelClient kestrelClient, String queueName) throws ParseError, IOException { String[] sentences = new String[] { "the cow jumped over the moon", "an apple a day keeps the doctor away", "four score and seven years ago", "snow white and the seven dwarfs", "i am at two with nature"}; Random _rand = new Random(); for(int i=1; i<=10; i++){ String sentence = sentences[_rand.nextInt(sentences.length)]; String val = "ID " + i + " " + sentence; boolean queueSucess = kestrelClient.queue(queueName, val); System.out.println("queueSucess=" +queueSucess+ " [" + val +"]"); } } 从 Kestrel 中移除数据 此方法从一个队列中取出一个数据,但并不把该数据从队列中删除: private static void dequeueItems(KestrelClient kestrelClient, String queueName) throws IOException, ParseError { for(int i=1; i<=12; i++){ Item item = kestrelClient.dequeue(queueName); if(item==null){ System.out.println("The queue (" + queueName + ") contains no items."); } else { byte[] data = item._data; String receivedVal = new String(data); System.out.println("receivedItem=" + receivedVal); } } 此方法会从队列中取出并移除数据: private static void dequeueAndRemoveItems(KestrelClient kestrelClient, String queueName) throws IOException, ParseError { for(int i=1; i<=12; i++){ Item item = kestrelClient.dequeue(queueName); if(item==null){ System.out.println("The queue (" + queueName + ") contains no items."); } else { int itemID = item._id; byte[] data = item._data; String receivedVal = new String(data); kestrelClient.ack(queueName, itemID); System.out.println("receivedItem=" + receivedVal); } } } 向 Kestrel 中连续添加数据 下面的程序可以向本地 Kestrel 服务的一个sentence_queue队列中连续添加句子,这也是我们的最后一个程序。 可以在命令行窗口中输入一个右中括号]并回车来停止程序。 import java.io.IOException; import java.io.InputStream; import java.util.Random; import backtype.storm.spout.KestrelClient; import backtype.storm.spout.KestrelClient.Item; import backtype.storm.spout.KestrelClient.ParseError; public class AddSentenceItemsToKestrel { /** * @param args */ public static void main(String[] args) { InputStream is = System.in; char closing_bracket = ']'; int val = closing_bracket; boolean aux = true; try { KestrelClient kestrelClient = null; String queueName = "sentence_queue"; while(aux){ kestrelClient = new KestrelClient("localhost",22133); queueSentenceItems(kestrelClient, queueName); kestrelClient.close(); Thread.sleep(1000); if(is.available()>0){ if(val==is.read()) aux=false; } } } catch (IOException e) { // TODO Auto-generated catch block e.printStackTrace(); } catch (ParseError e) { // TODO Auto-generated catch block e.printStackTrace(); } catch (InterruptedException e) { // TODO Auto-generated catch block e.printStackTrace(); } System.out.println("end"); } } 使用 KestrelSpout 下面的拓扑使用KestrelSpout从一个 Kestrel 队列中读取句子,并将句子分割成若干个单词(Bolt:SplitSentence),然后输出每个单词出现的次数(Bolt:WordCount)。数据处理的细节可以参考消息的可靠性保证一文。 TopologyBuilder builder = new TopologyBuilder(); builder.setSpout("sentences", new KestrelSpout("localhost",22133,"sentence_queue",new StringScheme())); builder.setBolt("split", new SplitSentence(), 10) .shuffleGrouping("sentences"); builder.setBolt("count", new WordCount(), 20) .fieldsGrouping("split", new Fields("word")); 运行 首先,以生产模式或者开发者模式启动你的本地 Kestrel 服务。 然后,等待大约 5 秒钟以防出现网络连接异常。 现在可以运行向队列中添加数据的程序,并启动 Storm 拓扑。程序启动的顺序并不重要。 如果你以 TOPOLOGY_DEBUG 模式运行拓扑你会观察到拓扑中 tuple 发送的细节信息。 转载自并发编程网 - ifeve.com

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

Apache Storm 官方文档 —— 配置开发环境

本文详细讲解了配置 Storm 开发环境的相关信息。简单地说,配置过程包含以下几个步骤: 下载Storm 发行版,将其解压缩并复制到你的PATH环境变量的bin目录中(也可以根据需要自定义安装目录 —— 译者注); 如果需要在远程集群中运行拓扑,则需要在~/.storm/storm.yaml文件中配置好集群的相关信息。 上述几步的详细内容如下。 什么是开发环境? Storm 包含两种操作模式:本地模式与远程模式(即集群模式 —— 译者注)。在本地模式下,你可以在本地机器上的一个进程中完成所有的开发、测试拓扑的工作。而在远程模式下,为了运行拓扑,你需要先向服务器集群提交该拓扑。 Storm 的开发环境已经为你准备好了一切,因此,你可以在本地模式下完成开发、测试拓扑的工作,将拓扑打包并提交到远程服务器,并在远程服务器集群上运行或者终止拓扑。 我们再来回顾一下本地机器与远程集群之间的关系。Storm 集群是由一个称为 “Nimbus” 的主节点管理的。本地机器通过与 Nimbus 通信来提交代码(代码已经打包为 jar 格式),这样代码文件中包含的拓扑就可以在集群中运行。Nimbus 会小心地维护着代码在集群中的分布式结构,并为待运行的拓扑分配 worker。本地机器可以使用一个称为storm的命令行客户端来与 Nimbus 进行通信。不过,storm客户端仅用于远程模式,不能用于本地模式下开发、测试拓扑。 在本地机器上安装 Storm 如果要从本地机器上直接向远程集群提交拓扑,你需要在本地机器上安装 Storm 程序。本地的 Storm 程序可以提供与远程集群交互的storm客户端。在安装本地 Storm 之前,你需要从这里下载一个 Storm 安装程序并将其解压到你的电脑的某个位置。然后将 Storm 的bin/目录添加到你的PATH环境变量中,确保bin/storm脚本可以直接运行。 在本地机器上安装的 Storm 仅能用于与远程集群的交互。对于本地模式下的开发、测试拓扑,推荐使用 Maven 来将 Storm 添加到你的项目的开发依赖中。关于 Maven 的使用请参考此文。 在远程集群上开始/终止拓扑的运行 在上一步中我们已经安装好了本地的storm客户端。接下来就需要告诉客户端需要连接哪一个 Storm 集群。这可以通过在~/.storm/storm.yaml文件中填写 Storm 集群的主节点的 host 地址来实现: nimbus.host: "123.45.678.890" 另外,如果你在 AWS 上应用storm-deploy项目来配置 Storm 集群,它会自动配置好你的~/.storm/storm.yaml文件。你也可以使用attach命令手动配置附属的 Storm 集群(或者在多个集群之间切换): lein run :deploy --attach --name mystormcluster 更多内容请参考 storm-deploy 项目的wiki。 转载自并发编程网 - ifeve.com

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

MaxCompute_2_MaxCompute数据迁移文档

免费开通大数据服务:https://www.aliyun.com/product/odps 乍一看标题会以为是不是作者写错了,怎么会有从MaxCompute到MaxCompute迁移数据的场景呢?在实际使用中已经有客户遇到了这种场景,比如:两个网络互通的专有云环境之间数据迁移、公共云数加DataIDE上两个云账号之间数据迁移、还有网络不通的两个MaxCompute项目数据迁移等等,下面我们逐个场景介绍。 场景一:两个网络互通的专有云MaxCompute环境之间数据迁移 这种场景需要先从源MaxCompute中导出元数据DDL,在目标MaxCompute中初始化表,然后借助DataX工具完成数据迁移,步骤如下: 1. 安装配置ODPS客户端 https://help.aliyun.com/document_detail/2

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

腾讯云软件源

腾讯云软件源

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

Nacos

Nacos

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

Sublime Text

Sublime Text

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

用户登录
用户注册