【问题标题】:Why does httpcomponents slow down my topology after the first processing of tuples?为什么在第一次处理元组后 httpcomponents 会减慢我的拓扑?
【发布时间】:2016-12-08 11:06:50
【问题描述】:

我已经构建了一个 Storm-topology,它通过 kafka-spout 从 Apache-Kafka 接收元组,将此数据(使用另一个螺栓)作为字符串写入本地系统上的 .txt 文件中,并从之后我的 PostBolt。

两个螺栓都连接到 Kafka-Spout。

如果我在没有 PostBolt 的情况下测试拓扑,一切正常。但是如果我将螺栓添加到拓扑中,整个拓扑会由于某种原因被阻塞。

有没有人遇到同样的问题或者有什么提示给我,是什么原因造成的?

我了解到存在一些 CloseableHttpClient 或 CloseableHttpResponse 阻塞线程无法工作的问题......在这种情况下可能是同样的问题吗?


我的 PostBolt 代码:

public class PostBolt extends BaseRichBolt {

private CloseableHttpClient httpclient; 

@Override
public final void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
    //empty for now
}

@Override
public final void execute(Tuple tuple) {

    //create HttpClient:
    httpclient = HttpClients.createDefault();
    String url = "http://xxx.xxx.xx.xxx:8080/HTTPServlet/httpservlet";
    HttpPost post = new HttpPost(url);

    post.setHeader("str1", "TEST TEST TEST");

    try {
        CloseableHttpResponse postResponse;
        postResponse = httpclient.execute(post);
        System.out.println(postResponse.getStatusLine());
        System.out.println("=====sending POST=====");
        HttpEntity postEntity = postResponse.getEntity();
        //do something useful with the response body
        //and ensure that it is fully consumed
        EntityUtils.consume(postEntity);
        postResponse.close();
    }catch (Exception e){
         e.printStackTrace();
    }
}

@Override
public void declareOutputFields(OutputFieldsDeclarer declarer) {
    declarer.declare(new Fields("HttpPost"));
}}

我的拓扑代码:

public static void main(String[] args) throws Exception {

    /**
    *   create a config for Kafka-Spout (and Kafka-Bolt)
    */
    Config config = new Config();
    config.setDebug(true);
    config.put(Config.TOPOLOGY_MAX_SPOUT_PENDING, 1);
    //setup zookeeper connection
    String zkConnString = "localhost:2181";
    //define Kafka topic for the spout
    String topic = "mytopic";
    //assign the zookeeper connection to brokerhosts
    BrokerHosts hosts = new ZkHosts(zkConnString);

    //setting up spout properties
    SpoutConfig kafkaSpoutConfig = new SpoutConfig(hosts, topic, "/" +topic, UUID.randomUUID().toString());
    kafkaSpoutConfig.bufferSizeBytes = 1024 * 1024 * 4;
    kafkaSpoutConfig.fetchSizeBytes = 1024 * 1024 * 4;
    kafkaSpoutConfig.scheme = new SchemeAsMultiScheme(new StringScheme());

    /**
    *   Build the Topology by linking the spout and bolts together
    */
    TopologyBuilder builder = new TopologyBuilder();
    builder.setSpout("kafka-spout", new KafkaSpout(kafkaSpoutConfig));
    builder.setBolt("printer-bolt", new PrinterBolt()).shuffleGrouping("kafka-spout");
    builder.setBolt("post-bolt", new PostBolt()).shuffleGrouping("kafka-spout");

    /**
    *   Check if we're running locally or on a real cluster
    */
    if (args != null && args.length >0) {
        config.setNumWorkers(6);
        config.setNumAckers(6);
        config.setMaxSpoutPending(100);
        config.setMessageTimeoutSecs(20);
        StormSubmitter.submitTopology("StormKafkaTopology", config, builder.createTopology());
    } else {
        config.setMaxTaskParallelism(3);
        config.setNumWorkers(6);
        LocalCluster cluster = new LocalCluster();
        cluster.submitTopology("StormKafkaTopology", config, builder.createTopology());
        //Utils.sleep(100000);
        //cluster.killTopology("StormKafkaTopology");
        //cluster.shutdown();
    }
}}

【问题讨论】:

    标签: java apache-kafka apache-storm apache-httpcomponents


    【解决方案1】:

    在我看来,您已经回答了您的问题,但是是的...根据this answer,您应该使用 PoolingHttpClientConnectionManager,因为您将在多线程环境中运行。

    编辑:

    public class PostBolt extends BaseRichBolt {
        private static Logger LOG = LoggerFactory.getLogger(PostBolt.class);
        private CloseableHttpClient httpclient;
        private OutputCollector _collector;        
    
        @Override
        public final void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
            httpclient = HttpClients.createDefault();
            _collector = collector;
        }
    
        @Override
        public final void execute(Tuple tuple) {
            String url = "http://xxx.xxx.xx.xxx:8080/HTTPServlet/httpservlet";
            HttpPost post = new HttpPost(url);
            post.setHeader("str1", "TEST TEST TEST");
    
            CloseableHttpResponse postResponse = httpclient.execute(post);
            try {
                LOG.info(postResponse.getStatusLine());
                LOG.info("=====sending POST=====");
                HttpEntity postEntity = postResponse.getEntity();
                //do something useful with the response body
                //and ensure that it is fully consumed
                EntityUtils.consume(postEntity);
                postResponse.close();
            }catch (Exception e){
               LOG.error("SolrIndexerBolt prepare error", e);
               _collector.reportError(e);
            } finally {
               postResponse.close()
            }
    
        }
    
        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("HttpPost"));
        }
    
    }
    

    【讨论】:

    • 感谢您的回答。可悲的是,我已经尝试使用 PoolingHttpClientConnectionManager,但是如果我添加 PostBolt,某些东西仍然会阻塞拓扑。您链接到我的帖子的 OP 也以另一种方式解决了他的问题:'我找到了阻塞的原因。由于来自服务器的 NOCONTENT 响应采用了正文,因此无法释放连接。这是测试的一个大错误,因为我只有一个环境来测试 NOCONTENT 响应。 – tank1920 2014 年 7 月 24 日 16:24 '
    • 好的,看起来应该可以了。我添加了一个代码示例供您尝试记录和关闭 postResponse。
    • 非常感谢!您的示例帮助我最大限度地减少拓扑中的滞后。但是,如果我同时使用 PostBolt 和 PrinterBolt...
    【解决方案2】:

    好的,我根据这条评论https://stackoverflow.com/a/32080845/7208987发现了问题

    Kafka Spout 将继续重新发送元组,这些元组没有被发送到的“端点”确认。

    所以我只需要确认螺栓内的传入元组,拓扑结构就没有了。

    (我发现了问题,因为printerbolt确实继续写入,即使kafkaspout没有进一步的输入)。

    【讨论】:

      猜你喜欢
      • 2018-08-08
      • 1970-01-01
      • 2017-11-14
      • 1970-01-01
      • 2023-03-20
      • 2018-06-25
      • 1970-01-01
      • 2021-08-01
      • 1970-01-01
      相关资源
      最近更新 更多