【问题标题】:Why are my storm topology not acking when i send tuple to ElasticSearch当我将元组发送到 ElasticSearch 时,为什么我的风暴拓扑不响应
【发布时间】:2020-05-10 21:27:29
【问题描述】:

我刚开始使用 Storm,我刚刚开始了数据架构师培训课程,正是在这种情况下,我面临着今天给你带来的问题。

我通过名为 CurrentPriceSpout 的 KafkaSpout 接收来自 kakfa 的消息。到目前为止,一切正常。然后,在我的 CurrentPriceBolt 中,我重新发布了一个元组,以便我的数据使用 EsCurrentPriceBolt 写入 ElasticSearch。问题就在这里。我无法将我的数据直接写入 ElasticSearch,只有在我删除拓扑时才会写入。

是否有一个 Storm 参数可以通过检索确认来强制写入元组?

我尝试添加选项“.addConfiguration(Config.TOPOLOGY_TICK_TUPLE_FREQ_SECS, 5)”,元组在 ElasticSearch 中写得很好,但没有得到承认。所以 Storm 会无限期地重写它们。

感谢您的帮助 蒂埃里

【问题讨论】:

  • 杀死拓扑做我想做的事,但我不能不断地杀死和添加拓扑杀死拓扑:1.刷新消息 2.发送确认是一个风暴命令或风暴选项,它们执行与杀死相同的事情拓扑
  • 您考虑过只使用 Kafka Connect 吗?这比你在这里做的更容易docs.confluent.io/current/connect/kafka-connect-elasticsearch/…
  • 如果我有选择,我可以这样做,是的!但这是一个评估训练成果的项目,它将 kafka -->storm --> elasticsearch 工作流强加给我。
  • 明白。祝你好运:-D

标签: maven elasticsearch apache-kafka apache-storm


【解决方案1】:

我设法找到了问题的答案。 主要问题是 ES 的设计目的不是摄取与研究项目中生成的数据一样少的数据。默认情况下,ES 以 1000 个条目为单位写入数据。在这个项目中,我每 30 秒生成一个数据,或者每 500 分钟(或 8h20)生成一批 1000 个数据。

所以我详细查看了我的拓扑结构并使用了以下选项:

  • es.batch.size.entries: 1
  • es.storm.bolt.flush.entries.size: 1
  • topology.producer.batch.size: 1
  • topology.transfer.batch.size: 1

现在是这样的:

...
...

public class App 
{
    ...    
    ...    

    public static void main( String[] args ) throws AlreadyAliveException, InvalidTopologyException, AuthorizationException
    {
        ...
        ...

        StormTopology topology  = topologyBuilder.createTopology();                 // je crée ma topologie Storm
        String topologyName     = properties.getProperty("storm.topology.name");    // je nomme ma topologie
        StormSubmitter.submitTopology(topologyName, getTopologyConfig(properties), topology);               // je démarre ma topologie sur mon cluster storm
        System.out.println( "Topology on remote cluster : Started!" );              
    }


    private static Config getTopologyConfig(Properties properties)
    {
        Config stormConfig = new Config();
        stormConfig.put("topology.workers",                 Integer.parseInt(properties.getProperty("topology.workers")));
        stormConfig.put("topology.enable.message.timeouts", Boolean.parseBoolean(properties.getProperty("topology.enable.message.timeouts")));
        stormConfig.put("topology.message.timeout.secs",    Integer.parseInt(properties.getProperty("topology.message.timeout.secs")));
        stormConfig.put("topology.transfer.batch.size",     Integer.parseInt(properties.getProperty("topology.transfer.batch.size")));
        stormConfig.put("topology.producer.batch.size",     Integer.parseInt(properties.getProperty("topology.producer.batch.size")));      
        return stormConfig;
    }

    ...    
    ...    
    ...    
}

现在它可以工作了!!!

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-08-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-08-21
    • 1970-01-01
    相关资源
    最近更新 更多