【问题标题】:How to specify zookeeper info while instantiating apache storm LocalCluster如何在实例化 apache Storm LocalCluster 时指定 Zookeeper 信息
【发布时间】:2017-03-28 17:52:16
【问题描述】:

我想实例化一个 LocalCluster() 并阻止它运行自己的嵌入式 zookeeper,并改用我的。

关于此问题:“https://issues.apache.org/jira/browse/STORM-213”已在 0.9.3 版本中解决。

我可以提供一个示例代码吗?

PS:我正在集成测试我的storm拓扑,我使用kafka和zookeeper作为storm的输入。 当我没有将 zookeeper 信息指定给 localcluster 时,我在“LocalCluster localCluster = new LocalCluster()”行收到此异常:

2016-06-08 12:16:56,785 WARN  [Thread-30] jmx.MBeanRegistry (MBeanRegistry.java:register(100)) - Failed to register MBean StandaloneServer_port-1
2016-06-08 12:16:56,785 WARN  [Thread-30] server.ZooKeeperServer (ZooKeeperServer.java:registerJMX(387)) - Failed to register with JMX
javax.management.InstanceAlreadyExistsException: org.apache.ZooKeeperService:name0=StandaloneServer_port-1
    at com.sun.jmx.mbeanserver.Repository.addMBean(Repository.java:437)
    at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerWithRepository(DefaultMBeanServerInterceptor.java:1898)
    at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerDynamicMBean(DefaultMBeanServerInterceptor.java:966)
    at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerObject(DefaultMBeanServerInterceptor.java:900)
    at com.sun.jmx.interceptor.DefaultMBeanServerInterceptor.registerMBean(DefaultMBeanServerInterceptor.java:324)
    at com.sun.jmx.mbeanserver.JmxMBeanServer.registerMBean(JmxMBeanServer.java:522)
    at org.apache.zookeeper.jmx.MBeanRegistry.register(MBeanRegistry.java:96)
    at org.apache.zookeeper.server.ZooKeeperServer.registerJMX(ZooKeeperServer.java:377)
    at org.apache.zookeeper.server.ZooKeeperServer.startup(ZooKeeperServer.java:410)
    at org.apache.zookeeper.server.NIOServerCnxnFactory.startup(NIOServerCnxnFactory.java:123)

当我为本地集群指定“storm.zookeeper.servers”和“storm.zookeeper.port”时,我在“localCluster.submitTopology()”行得到以下异常:

EndOfStreamException: Unable to read additional data from client sessionid 0x1552f0890b70000, likely client has closed socket
    at org.apache.zookeeper.server.NIOServerCnxn.doIO(NIOServerCnxn.java:228)
    at org.apache.zookeeper.server.NIOServerCnxnFactory.run(NIOServerCnxnFactory.java:208)
    at java.lang.Thread.run(Thread.java:745)

java.lang.NullPointerException
    at clojure.lang.Reflector.invokeInstanceMethod(Reflector.java:26)
    at org.apache.storm.testing$submit_local_topology.invoke(testing.clj:301)
    at org.apache.storm.LocalCluster$_submitTopology.invoke(LocalCluster.clj:49)
    at org.apache.storm.LocalCluster.submitTopology(Unknown Source)

【问题讨论】:

  • 我认为,在您交给 LocalCluster 以使其工作的配置中指定“storm.zookeeper.servers”属性就足够了。
  • tanx 为您的回复,我做到了,还设置了“storm.zookeeper.port”,但是当我提交拓扑时,我得到一个 NullPointerException!没有更多细节!
  • 你能分享堆栈跟踪吗?

标签: local apache-storm apache-zookeeper


【解决方案1】:

您可以使用重载LocalCluster 构造函数。

LocalCluster cluster = new LocalCluster("localhost", 2181L);

【讨论】:

    【解决方案2】:

    通过使用此链接http://storm.apache.org/releases/1.0.2/storm-kafka.html提供的配置,我能够使用 Kafka 作为 Storm 拓扑的输入

    我为storm config创建了一个特定的java类

    public class StormConfig {
    
    private String zooKeeperConnect;
    
    public StormConfig() {
    
    }
    public KafkaSpout getkafkaSpout(String topic){
        return new KafkaSpout(this.getSpoutConfig(topic));
    }
    
    public SpoutConfig getSpoutConfig(String topic) {
        SpoutConfig spoutConfig=new SpoutConfig(this.getZkHosts(), topic, "", topic);
        spoutConfig.scheme = new SchemeAsMultiScheme(new StringScheme());
        spoutConfig.startOffsetTime=kafka.api.OffsetRequest.EarliestTime();
        return spoutConfig;
    }
    
    public ZkHosts getZkHosts() {
        return new ZkHosts(getZooKeeperConnect());
    }
    
    public String getZooKeeperConnect() {
        return zooKeeperConnect;
    }
    
    public void setZooKeeperConnect(String zooKeeperConnect) {
        this.zooKeeperConnect = zooKeeperConnect;
    }
    }
    

    然后我在创建拓扑的初始 spout 时使用了此类中的方法:

    builder.setSpout("kafkaSpoutName", stormConfig.getkafkaSpout("topicName"))
                .setNumTasks(Constants.SPOUT_NUM_TASKS);
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-11-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多