小编Ris*_*abh的帖子

使用Java代码进行风暴拓扑重新平衡

我正在尝试重新平衡使用KafkaSpout的Storm拓扑.我的代码是:

    TopologyBuilder builder = new TopologyBuilder();
    Properties kafkaProps = new Properties();
    kafkaProps.put("zk.connect", "localhost:2181");
    kafkaProps.put("zk.connectiontimeout.ms", "1000000");
    kafkaProps.put("groupid", "storm");

    builder.setSpout( "kafkaSpout" , new KafkaSpout(kafkaProps, "test"), 3);
    builder.setBolt( "eventBolt", new EventBolt(), 2 ).shuffleGrouping( "kafkaSpout", "eventStream" );
    builder.setBolt( "tableBolt", new TableBolt(), 2 ).shuffleGrouping( "kafkaSpout", "tableStream");

    Map<String, Object> conf = new HashMap<String, Object>();
    conf.put(Config.TOPOLOGY_DEBUG, true);

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

    Utils.sleep( 1000*5 );

    List<TopologySummary> topologySummaries = cluster.getClusterInfo().get_topologies();
    for ( TopologySummary summary : topologySummaries ) {
        StormTopology topology = cluster.getTopology( summary.get_id() );
        RebalanceOptions …
Run Code Online (Sandbox Code Playgroud)

java topology apache-kafka apache-storm apache-zookeeper

6
推荐指数
1
解决办法
1722
查看次数

Kafka日志中的偏移量丢失 - 简单消费者无法继续

我有一个3节点kafka群集设置.我正在使用风暴来阅读来自kafka的消息.我系统中的每个主题都有7个分区.

现在我面临一个奇怪的问题.直到3天前,一切都运转良好.但是,现在看来我的风暴拓扑无法从2个分区 - #1和#4中专门读取.

我试着深入研究问题并发现在我的kafka日志中,对于这两个分区,缺少一个偏移,即在5964511之后,下一个偏移是5964513而不是5964512.

由于缺少偏移,Simple Consumer无法继续下一个偏移.我做错了什么或者它是一个已知的错误?

可能是这种行为的原因是什么?

我使用以下代码来读取有效偏移的窗口:

public static long getLastOffset(SimpleConsumer consumer, String topic, int partition,
                                 long whichTime, String clientName) {
    TopicAndPartition topicAndPartition = new TopicAndPartition(topic, partition);
    Map<TopicAndPartition, PartitionOffsetRequestInfo> requestInfoMap = new HashMap<TopicAndPartition, PartitionOffsetRequestInfo>();
    requestInfoMap.put(topicAndPartition, new PartitionOffsetRequestInfo(kafka.api.OffsetRequest.LatestTime(), 100));
    OffsetRequest request = new OffsetRequest( requestInfoMap, kafka.api.OffsetRequest.CurrentVersion() , clientName);
    OffsetResponse response = consumer.getOffsetsBefore(request);
    long[] validOffsets = response.offsets(topic, partition);
    for (long validOffset : validOffsets) {
        System.out.println(validOffset + " : ");
    }
    long largestOffset = validOffsets[0];
    long smallestOffset = validOffsets[validOffsets.length - 1]; …
Run Code Online (Sandbox Code Playgroud)

java apache-kafka apache-storm

5
推荐指数
1
解决办法
1231
查看次数