[
https://issues.apache.org/jira/browse/FLINK-3341?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=15137358#comment-15137358
]
ASF GitHub Bot commented on FLINK-3341:
---------------------------------------
Github user rmetzger commented on a diff in the pull request:
https://github.com/apache/flink/pull/1597#discussion_r52202024
--- 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 --
No. Check the new consumer config of Kafka 0.9:
http://kafka.apache.org/documentation.html#newconsumerconfigs
> 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)