gpt4 book ai didi

storm.kafka.ZkHosts.()方法的使用及代码示例

转载 作者:知者 更新时间:2024-03-15 10:36:49 27 4
gpt4 key购买 nike

本文整理了Java中storm.kafka.ZkHosts.<init>()方法的一些代码示例,展示了ZkHosts.<init>()的具体用法。这些代码示例主要来源于Github/Stackoverflow/Maven等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。ZkHosts.<init>()方法的具体详情如下:
包路径:storm.kafka.ZkHosts
类名称:ZkHosts
方法名:<init>

ZkHosts.<init>介绍

暂无

代码示例

代码示例来源:origin: qidasheng/storm-kafka-xlog

public XlogKafkaSpoutTopology(String kafkaZookeeper) {
  brokerHosts = new ZkHosts(kafkaZookeeper);
}

代码示例来源:origin: wurstmeister/storm-kafka-0.8-plus-test

public KafkaSpoutTestTopology(String kafkaZookeeper) {
  brokerHosts = new ZkHosts(kafkaZookeeper);
}

代码示例来源:origin: wurstmeister/storm-kafka-0.8-plus-test

public SentenceAggregationTopology(String kafkaZookeeper) {
  brokerHosts = new ZkHosts(kafkaZookeeper);
}

代码示例来源:origin: stackoverflow.com

import storm.kafka

TridentTopology topology = new TridentTopology();
BrokerHosts zk = new ZkHosts("localhost");
TridentKafkaConfig spoutConf = new TridentKafkaConfig(zk, "test-topic");
spoutConf.scheme = new SchemeAsMultiScheme(new StringScheme());
OpaqueTridentKafkaSpout spout = new OpaqueTridentKafkaSpout(spoutConf);

代码示例来源:origin: org.apache.eagle/eagle-metric-collection

hosts = new ZkHosts(zkConnString);
} else {
  hosts = new ZkHosts(zkConnString, brokerZkPath);

代码示例来源:origin: eshioji/trident-tutorial

public static void main(String[] args) throws Exception {
  Config conf = new Config();
  if (args.length == 2) {
    // Ready & submit the topology
    String name = args[0];
    BrokerHosts hosts = new ZkHosts(args[1]);
    TransactionalTridentKafkaSpout kafkaSpout = TestUtils.testTweetSpout(hosts);
    StormSubmitter.submitTopology(name, conf, buildTopology(kafkaSpout));
  }else{
    System.err.println("<topologyName> <zookeeperHost>");
  }
}

代码示例来源:origin: eshioji/trident-tutorial

public static void main(String[] args) throws Exception {
  Config conf = new Config();
  if (args.length == 2) {
    // Ready & submit the topology
    String name = args[0];
    BrokerHosts hosts = new ZkHosts(args[1]);
    TransactionalTridentKafkaSpout kafkaSpout = TestUtils.testTweetSpout(hosts);
    StormSubmitter.submitTopology(name, conf, buildTopology(kafkaSpout));
  }else{
    System.err.println("<topologyName> <zookeeperHost>");
  }
}

代码示例来源:origin: eshioji/trident-tutorial

public static void main(String[] args) throws Exception {
  Config conf = new Config();
  if (args.length == 2) {
    // Ready & submit the topology
    String name = args[0];
    BrokerHosts hosts = new ZkHosts(args[1]);
    TransactionalTridentKafkaSpout kafkaSpout = TestUtils.testTweetSpout(hosts);
    StormSubmitter.submitTopology(name, conf, buildTopology(kafkaSpout));
  }else{
    System.err.println("<topologyName> <zookeeperHost>");
  }
}

代码示例来源:origin: eshioji/trident-tutorial

public static void main(String[] args) throws Exception {
  Config conf = new Config();
  if (args.length == 2) {
    // Ready & submit the topology
    String name = args[0];
    BrokerHosts hosts = new ZkHosts(args[1]);
    TransactionalTridentKafkaSpout kafkaSpout = TestUtils.testTweetSpout(hosts);
    StormSubmitter.submitTopology(name, conf, buildTopology(kafkaSpout));
  }else{
    System.err.println("<topologyName> <zookeeperHost>");
  }
}

代码示例来源:origin: eshioji/trident-tutorial

public static void main(String[] args) throws Exception {
  Config conf = new Config();
  if (args.length == 2) {
    // Ready & submit the topology
    String name = args[0];
    BrokerHosts hosts = new ZkHosts(args[1]);
    TransactionalTridentKafkaSpout kafkaSpout = TestUtils.testTweetSpout(hosts);
    StormSubmitter.submitTopology(name, conf, buildTopology(kafkaSpout));
  }else{
    System.err.println("<topologyName> <zookeeperHost>");
  }
}

代码示例来源:origin: eshioji/trident-tutorial

public static void main(String[] args) throws Exception {
  Config conf = new Config();
  if (args.length == 2) {
    // Ready & submit the topology
    String name = args[0];
    BrokerHosts hosts = new ZkHosts(args[1]);
    TransactionalTridentKafkaSpout kafkaSpout = TestUtils.testTweetSpout(hosts);
    StormSubmitter.submitTopology(name, conf, buildTopology(kafkaSpout));
  }else{
    System.err.println("<topologyName> <zookeeperHost>");
  }
}

代码示例来源:origin: OpenSOC/opensoc-streaming

private boolean initializeKafkaSpout(String name) {
  try {
    BrokerHosts zk = new ZkHosts(config.getString("kafka.zk"));
    String input_topic = config.getString("spout.kafka.topic");
    SpoutConfig kafkaConfig = new SpoutConfig(zk, input_topic, "",
        input_topic);
    kafkaConfig.scheme = new SchemeAsMultiScheme(new RawScheme());
    kafkaConfig.forceFromStart = Boolean.valueOf("True");
    kafkaConfig.startOffsetTime = -1;
    builder.setSpout(name, new KafkaSpout(kafkaConfig),
        config.getInt("spout.kafka.parallelism.hint")).setNumTasks(
        config.getInt("spout.kafka.num.tasks"));
  } catch (Exception e) {
    e.printStackTrace();
    System.exit(0);
  }
  return true;
}

代码示例来源:origin: eshioji/trident-tutorial

public static void main(String[] args) throws Exception {
  Config conf = new Config();
  conf.setNumWorkers(6);
  if (args.length == 2) {
    // Ready & submit the topology
    String name = args[0];
    BrokerHosts hosts = new ZkHosts(args[1]);
    TransactionalTridentKafkaSpout kafkaSpout = TestUtils.testTweetSpout(hosts);
    StormSubmitter.submitTopology(name, conf, buildTopology(kafkaSpout));
  }else{
    System.err.println("<topologyName> <zookeeperHost>");
  }
}

代码示例来源:origin: XavientInformationSystems/Data-Ingestion-Platform

public static KafkaSpout getKafkaSpout(String topic, String zkHost, String zkPort, Boolean rewind) {
  String zkRoot = "/" + topic;
  String zkSpoutId = topic;
  BrokerHosts hosts = new ZkHosts(zkHost + ":" + zkPort);
  SpoutConfig spoutCfg = new SpoutConfig(hosts, topic, zkRoot, zkSpoutId);
  if (rewind) {
    spoutCfg.ignoreZkOffsets = true;
    //spoutCfg.startOffsetTime = kafka.api.OffsetRequest.EarliestTime();
  }
  spoutCfg.scheme = new SchemeAsMultiScheme(new StringScheme());
  KafkaSpout kafkaSpout = new KafkaSpout(spoutCfg);
  return kafkaSpout;
}

代码示例来源:origin: stackoverflow.com

ZkHosts zkHosts = new ZkHosts(zookeeperHost);

代码示例来源:origin: Big-Data-Manning/big-data-code

public static LocalCluster run() throws Exception {
    TopologyBuilder builder = new TopologyBuilder();
    SpoutConfig spoutConfig = new SpoutConfig(
        new ZkHosts("127.0.0.1:2181"),
        "pageviews",
        "/kafkastorm",
        "uniquesSpeedLayer"
    );
    spoutConfig.scheme = new PageviewScheme();

    builder.setSpout("pageviews",
        new KafkaSpout(spoutConfig), 2);

    builder.setBolt("extract-filter",
        new ExtractFilterBolt(), 4)
        .shuffleGrouping("pageviews");
    builder.setBolt("cassandra",
        new UpdateCassandraBolt(), 4)
        .fieldsGrouping("extract-filter",
            new Fields("domain"));

    LocalCluster cluster = new LocalCluster();
    Config conf = new Config();
    cluster.submitTopology("uniques", conf, builder.createTopology());

    return cluster;
  }
}

代码示例来源:origin: mayconbordin/storm-applications

@Override
protected void initialize() {
  String parserClass = config.getString(getConfigKey(BaseConf.SPOUT_PARSER));
  String host        = config.getString(getConfigKey(BaseConf.KAFKA_HOST));
  String topic       = config.getString(getConfigKey(BaseConf.KAFKA_SPOUT_TOPIC));
  String consumerId  = config.getString(getConfigKey(BaseConf.KAFKA_CONSUMER_ID));
  String path        = config.getString(getConfigKey(BaseConf.KAFKA_ZOOKEEPER_PATH));
  
  Parser parser = (Parser) ClassLoaderUtils.newInstance(parserClass, "parser", LOG);
  parser.initialize(config);
  
  Fields defaultFields = fields.get(Utils.DEFAULT_STREAM_ID);
  if (defaultFields == null) {
    throw new RuntimeException("A KafkaSpout must have a default stream");
  }
  
  brokerHosts = new ZkHosts(host);
  
  SpoutConfig spoutConfig = new SpoutConfig(brokerHosts, topic, path, consumerId);
  spoutConfig.scheme = new SchemeAsMultiScheme(new ParserScheme(parser, defaultFields));
  
  spout = new storm.kafka.KafkaSpout(spoutConfig);
  spout.open(config, context, collector);
}

代码示例来源:origin: Big-Data-Manning/big-data-code

TridentKafkaConfig kafkaConfig =
    new TridentKafkaConfig(
        new ZkHosts(
            "127.0.0.1:2181"),
        "pageviews"

代码示例来源:origin: Big-Data-Manning/big-data-code

TridentKafkaConfig kafkaConfig =
    new TridentKafkaConfig(
        new ZkHosts(
            "127.0.0.1:2181"),
        "pageviews"

代码示例来源:origin: jasonTangxd/recommendSys

ZkHosts zkHosts = new ZkHosts(brokerZkStr);
String topic = "recommender";

27 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com