You are viewing a plain text version of this content. The canonical link for it is here.
Posted to issues@storm.apache.org by "Janith Kaiprath Valiyalappil (JIRA)" <ji...@apache.org> on 2017/09/05 05:11:00 UTC

[jira] [Created] (STORM-2721) CLONE - Add timestamp based FirstPollOffsetStrategy in KafkaTridentSpoutOpaque

Janith Kaiprath Valiyalappil created STORM-2721:
---------------------------------------------------

             Summary: CLONE - Add timestamp based FirstPollOffsetStrategy in KafkaTridentSpoutOpaque
                 Key: STORM-2721
                 URL: https://issues.apache.org/jira/browse/STORM-2721
             Project: Apache Storm
          Issue Type: Improvement
          Components: storm-kafka-client, trident
    Affects Versions: 1.1.1
            Reporter: Janith Kaiprath Valiyalappil
            Priority: Minor


Offsets for a given partition at a particular timestamp can now be found using offsetsForTimes API. https://kafka.apache.org/0110/javadoc/org/apache/kafka/clients/consumer/KafkaConsumer.html#offsetsForTimes(java.util.Map).

One way to make use of this api would be to :
Add a new option for FirstPollOffsetStrategy called TIMESTAMP 
Add a new startTimeStamp option to KafkaSpoutConfig, which would be used only when FirstPollOffsetStrategy is set to TIMESTAMP.

Later in the KafkaTridentSpoutEmitter, when we do the first seek, we can do something like :

{code}
            if(firstPollOffsetStrategy.equals(TIMESTAMP)) {
                try {
                    startTimeStampOffset =
                        kafkaConsumer.offsetsForTimes(Collections.singletonMap(tp, startTimeStamp)).get(tp).offset();
                } catch (IllegalArgumentException e) {
                    LOG.error("Illegal timestamp {} provided for tp {} ",startTimeStamp,tp.toString());
                } catch (UnsupportedVersionException e) {
                    LOG.error("Kafka Server do not support offsetsForTimes(), probably < 0.10.1",e);
                }

                if(startTimeStampOffset!=null) {
                    LOG.info("Kafka consumer offset reset for TopicPartition {}, TimeStamp {}, Offset {}",tp,startTimeStamp,startTimeStampOffset);
                    kafkaConsumer.seek(tp, startTimeStampOffset);
                } else {
                    LOG.info("Kafka consumer offset reset by timestamp failed for TopicPartition {}, TimeStamp {}, Offset {}. Restart with a different Strategy ",tp,startTimeStamp,startTimeStampOffset);
                }
            }
{code}





--
This message was sent by Atlassian JIRA
(v6.4.14#64029)