[ 
https://issues.apache.org/jira/browse/FLINK-3341?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=15137347#comment-15137347
 ] 

ASF GitHub Bot commented on FLINK-3341:
---------------------------------------

Github user shikhar commented on a diff in the pull request:

    https://github.com/apache/flink/pull/1597#discussion_r52201602
  
    --- Diff: 
flink-streaming-connectors/flink-connector-kafka-0.8/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/LegacyFetcher.java
 ---
    @@ -576,7 +576,8 @@ private static void getLastOffset(SimpleConsumer 
consumer, List<FetchPartition>
     
                private static long getInvalidOffsetBehavior(Properties config) 
{
                        long timeType;
    -                   if 
(config.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, 
"latest").equals("latest")) {
    +                   String val = 
config.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "largest");
    +                   if (val.equals("largest") || val.equals("latest")) { // 
largest is kafka 0.8, latest is kafka 0.9
    --- End diff --
    
    > latest is kafka 0.9
    
    I think 'latest' is just flink :)


> Kafka connector's 'auto.offset.reset' inconsistent with Kafka
> -------------------------------------------------------------
>
>                 Key: FLINK-3341
>                 URL: https://issues.apache.org/jira/browse/FLINK-3341
>             Project: Flink
>          Issue Type: Bug
>            Reporter: Shikhar Bhushan
>            Assignee: Robert Metzger
>            Priority: Minor
>
> Kafka docs talk of valid "auto.offset.reset" values being "smallest" or 
> "largest"
> https://kafka.apache.org/08/configuration.html
> The {{LegacyFetcher}} looks for "latest" and otherwise defaults to "smallest"
> cc [~rmetzger]



--
This message was sent by Atlassian JIRA
(v6.3.4#6332)

Reply via email to