Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,6 @@
import java.util.Iterator;
import java.util.Map;
import java.util.Optional;
import java.util.OptionalInt;
import java.util.OptionalLong;
import java.util.concurrent.CountDownLatch;
import java.util.regex.Pattern;
import java.util.Collections;
Expand Down Expand Up @@ -68,11 +66,8 @@ public static void main(String[] args) throws Exception {

public static void run(ConsoleConsumerOptions opts) {
messageCount = 0;
long timeoutMs = opts.timeoutMs() >= 0 ? opts.timeoutMs() : Long.MAX_VALUE;
Consumer<byte[], byte[]> consumer = new KafkaConsumer<>(opts.consumerProps(), new ByteArrayDeserializer(), new ByteArrayDeserializer());
ConsumerWrapper consumerWrapper = opts.partitionArg().isPresent()
? new ConsumerWrapper(Optional.of(opts.topicArg()), opts.partitionArg(), OptionalLong.of(opts.offsetArg()), Optional.empty(), consumer, timeoutMs)
: new ConsumerWrapper(Optional.of(opts.topicArg()), OptionalInt.empty(), OptionalLong.empty(), Optional.ofNullable(opts.includedTopicsArg()), consumer, timeoutMs);
ConsumerWrapper consumerWrapper = new ConsumerWrapper(opts, consumer);

addShutdownHook(consumerWrapper, opts);

Expand Down Expand Up @@ -148,43 +143,25 @@ static boolean checkErr(PrintStream output) {
}

public static class ConsumerWrapper {
final Optional<String> topic;
final OptionalInt partitionId;
final OptionalLong offset;
final Optional<String> includedTopics;
final Consumer<byte[], byte[]> consumer;
final long timeoutMs;
final Time time = Time.SYSTEM;
final long timeoutMs;
final Consumer<byte[], byte[]> consumer;

Iterator<ConsumerRecord<byte[], byte[]>> recordIter = Collections.emptyIterator();

public ConsumerWrapper(Optional<String> topic,
OptionalInt partitionId,
OptionalLong offset,
Optional<String> includedTopics,
Consumer<byte[], byte[]> consumer,
long timeoutMs) {
this.topic = topic;
this.partitionId = partitionId;
this.offset = offset;
this.includedTopics = includedTopics;
public ConsumerWrapper(ConsoleConsumerOptions opts, Consumer<byte[], byte[]> consumer) {
this.consumer = consumer;
this.timeoutMs = timeoutMs;

if (topic.isPresent() && partitionId.isPresent() && offset.isPresent() && !includedTopics.isPresent()) {
seek(topic.get(), partitionId.getAsInt(), offset.getAsLong());
} else if (topic.isPresent() && partitionId.isPresent() && !offset.isPresent() && !includedTopics.isPresent()) {
// default to latest if no offset is provided
seek(topic.get(), partitionId.getAsInt(), ListOffsetsRequest.LATEST_TIMESTAMP);
} else if (topic.isPresent() && !partitionId.isPresent() && !offset.isPresent() && !includedTopics.isPresent()) {
consumer.subscribe(Collections.singletonList(topic.get()));
} else if (!topic.isPresent() && !partitionId.isPresent() && !offset.isPresent() && includedTopics.isPresent()) {
consumer.subscribe(Pattern.compile(includedTopics.get()));
timeoutMs = opts.timeoutMs();
Optional<String> topic = opts.topicArg();

if (topic.isPresent()) {
if (opts.partitionArg().isPresent()) {
seek(topic.get(), opts.partitionArg().getAsInt(), opts.offsetArg());
} else {
consumer.subscribe(Collections.singletonList(topic.get()));
}
} else {
throw new IllegalArgumentException("An invalid combination of arguments is provided. " +
Comment thread
wernerdv marked this conversation as resolved.
"Exactly one of 'topic' or 'include' must be provided. " +
"If 'topic' is provided, an optional 'partition' may also be provided. " +
"If 'partition' is provided, an optional 'offset' may also be provided, otherwise, consumption starts from the end of the partition.");
opts.includedTopicsArg().ifPresent(topics -> consumer.subscribe(Pattern.compile(topics)));
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.OptionalInt;
import java.util.Properties;
import java.util.Random;
Expand All @@ -55,7 +55,7 @@ public final class ConsoleConsumerOptions extends CommandDefaultOptions {
private final OptionSpec<String> messageFormatterConfigOpt;
private final OptionSpec<?> resetBeginningOpt;
private final OptionSpec<Integer> maxMessagesOpt;
private final OptionSpec<Integer> timeoutMsOpt;
private final OptionSpec<Long> timeoutMsOpt;
private final OptionSpec<?> skipMessageOnErrorOpt;
private final OptionSpec<String> bootstrapServerOpt;
private final OptionSpec<String> keyDeserializerOpt;
Expand All @@ -66,6 +66,7 @@ public final class ConsoleConsumerOptions extends CommandDefaultOptions {

private final Properties consumerProps;
private final long offset;
private final long timeoutMs;
private final MessageFormatter formatter;

public ConsoleConsumerOptions(String[] args) throws IOException {
Expand Down Expand Up @@ -139,7 +140,7 @@ public ConsoleConsumerOptions(String[] args) throws IOException {
timeoutMsOpt = parser.accepts("timeout-ms", "If specified, exit if no message is available for consumption for the specified interval.")
.withRequiredArg()
.describedAs("timeout_ms")
.ofType(Integer.class);
.ofType(Long.class);
skipMessageOnErrorOpt = parser.accepts("skip-message-on-error", "If there is an error when processing a message, " +
"skip it instead of halt.");
bootstrapServerOpt = parser.accepts("bootstrap-server", "REQUIRED: The server(s) to connect to.")
Expand Down Expand Up @@ -184,12 +185,13 @@ public ConsoleConsumerOptions(String[] args) throws IOException {
Set<String> groupIdsProvided = checkConsumerGroup(consumerPropsFromFile, extraConsumerProps);
consumerProps = buildConsumerProps(consumerPropsFromFile, extraConsumerProps, groupIdsProvided);
offset = parseOffset();
timeoutMs = parseTimeoutMs();
formatter = buildFormatter();
}

private void checkRequiredArgs() {
List<String> topicOrFilterArgs = new ArrayList<>(Arrays.asList(topicArg(), includedTopicsArg()));
topicOrFilterArgs.removeIf(Objects::isNull);
List<Optional<String>> topicOrFilterArgs = new ArrayList<>(Arrays.asList(topicArg(), includedTopicsArg()));
topicOrFilterArgs.removeIf(arg -> !arg.isPresent());
// user need to specify value for either --topic or one of the include filters options (--include or --whitelist)
if (topicOrFilterArgs.size() != 1) {
CommandLineUtils.printUsageAndExit(parser, "Exactly one of --include/--topic is required. " +
Expand Down Expand Up @@ -322,6 +324,11 @@ private void invalidOffset(String offset) {
"'earliest', 'latest', or a non-negative long.");
}

private long parseTimeoutMs() {
long timeout = options.has(timeoutMsOpt) ? options.valueOf(timeoutMsOpt) : -1;
return timeout >= 0 ? timeout : Long.MAX_VALUE;
}

private MessageFormatter buildFormatter() {
MessageFormatter formatter = null;
try {
Expand Down Expand Up @@ -365,16 +372,16 @@ OptionalInt partitionArg() {
return OptionalInt.empty();
}

String topicArg() {
return options.valueOf(topicOpt);
Optional<String> topicArg() {
return options.has(topicOpt) ? Optional.of(options.valueOf(topicOpt)) : Optional.empty();
}

int maxMessages() {
return options.has(maxMessagesOpt) ? options.valueOf(maxMessagesOpt) : -1;
}

int timeoutMs() {
return options.has(timeoutMsOpt) ? options.valueOf(timeoutMsOpt) : -1;
long timeoutMs() {
return timeoutMs;
}

boolean enableSystestEventsLogging() {
Expand All @@ -385,10 +392,10 @@ String bootstrapServer() {
return options.valueOf(bootstrapServerOpt);
}

String includedTopicsArg() {
Optional<String> includedTopicsArg() {
return options.has(includeOpt)
? options.valueOf(includeOpt)
: options.valueOf(whitelistOpt);
? Optional.of(options.valueOf(includeOpt))
: Optional.ofNullable(options.valueOf(whitelistOpt));
}

Properties formatterArgs() throws IOException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,12 +48,12 @@ public void shouldParseValidConsumerValidConfig() throws IOException {
ConsoleConsumerOptions config = new ConsoleConsumerOptions(args);

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertTrue(config.fromBeginning());
assertFalse(config.enableSystestEventsLogging());
assertFalse(config.skipMessageOnError());
assertEquals(-1, config.maxMessages());
assertEquals(-1, config.timeoutMs());
assertEquals(Long.MAX_VALUE, config.timeoutMs());
}

@Test
Expand All @@ -67,7 +67,7 @@ public void shouldParseIncludeArgument() throws IOException {
ConsoleConsumerOptions config = new ConsoleConsumerOptions(args);

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("includeTest*", config.includedTopicsArg());
assertEquals("includeTest*", config.includedTopicsArg().orElse(""));
assertTrue(config.fromBeginning());
}

Expand All @@ -82,7 +82,7 @@ public void shouldParseWhitelistArgument() throws IOException {
ConsoleConsumerOptions config = new ConsoleConsumerOptions(args);

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("whitelistTest*", config.includedTopicsArg());
assertEquals("whitelistTest*", config.includedTopicsArg().orElse(""));
assertTrue(config.fromBeginning());
}

Expand All @@ -96,7 +96,7 @@ public void shouldIgnoreWhitelistArgumentIfIncludeSpecified() throws IOException
};
ConsoleConsumerOptions config = new ConsoleConsumerOptions(args);
assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("includeTest*", config.includedTopicsArg());
assertEquals("includeTest*", config.includedTopicsArg().orElse(""));
assertTrue(config.fromBeginning());
}

Expand All @@ -112,7 +112,7 @@ public void shouldParseValidSimpleConsumerValidConfigWithNumericOffset() throws
ConsoleConsumerOptions config = new ConsoleConsumerOptions(args);

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertTrue(config.partitionArg().isPresent());
assertEquals(0, config.partitionArg().getAsInt());
assertEquals(3, config.offsetArg());
Expand Down Expand Up @@ -191,7 +191,7 @@ public void shouldParseValidSimpleConsumerValidConfigWithStringOffset() throws E
ConsoleConsumerOptions config = new ConsoleConsumerOptions(args);

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertTrue(config.partitionArg().isPresent());
assertEquals(0, config.partitionArg().getAsInt());
assertEquals(-1, config.offsetArg());
Expand All @@ -211,7 +211,7 @@ public void shouldParseValidConsumerConfigWithAutoOffsetResetLatest() throws IOE
Properties consumerProperties = config.consumerProps();

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertFalse(config.fromBeginning());
assertEquals("latest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG));
}
Expand All @@ -228,7 +228,7 @@ public void shouldParseValidConsumerConfigWithAutoOffsetResetEarliest() throws I
Properties consumerProperties = config.consumerProps();

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertFalse(config.fromBeginning());
assertEquals("earliest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG));
}
Expand All @@ -246,7 +246,7 @@ public void shouldParseValidConsumerConfigWithAutoOffsetResetAndMatchingFromBegi
Properties consumerProperties = config.consumerProps();

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertTrue(config.fromBeginning());
assertEquals("earliest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG));
}
Expand All @@ -262,7 +262,7 @@ public void shouldParseValidConsumerConfigWithNoOffsetReset() throws IOException
Properties consumerProperties = config.consumerProps();

assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertFalse(config.fromBeginning());
assertEquals("latest", consumerProperties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG));
}
Expand Down Expand Up @@ -442,7 +442,7 @@ public void shouldParseGroupIdFromBeginningGivenTogether() throws IOException {

ConsoleConsumerOptions config = new ConsoleConsumerOptions(args);
assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertEquals(-2, config.offsetArg());
assertTrue(config.fromBeginning());

Expand All @@ -455,7 +455,7 @@ public void shouldParseGroupIdFromBeginningGivenTogether() throws IOException {

config = new ConsoleConsumerOptions(args);
assertEquals("localhost:9092", config.bootstrapServer());
assertEquals("test", config.topicArg());
assertEquals("test", config.topicArg().orElse(""));
assertEquals(-1, config.offsetArg());
assertFalse(config.fromBeginning());
}
Expand Down Expand Up @@ -618,4 +618,30 @@ public void testParseOffset() throws Exception {
Exit.resetExitProcedure();
}
}

@Test
public void testParseTimeoutMs() throws Exception {
String[] withoutTimeoutMs = new String[]{
"--bootstrap-server", "localhost:9092",
"--topic", "test",
"--partition", "0"
};
assertEquals(Long.MAX_VALUE, new ConsoleConsumerOptions(withoutTimeoutMs).timeoutMs());

String[] negativeTimeoutMs = new String[]{
"--bootstrap-server", "localhost:9092",
"--topic", "test",
"--partition", "0",
"--timeout-ms", "-100"
};
assertEquals(Long.MAX_VALUE, new ConsoleConsumerOptions(negativeTimeoutMs).timeoutMs());

String[] validTimeoutMs = new String[]{
"--bootstrap-server", "localhost:9092",
"--topic", "test",
"--partition", "0",
"--timeout-ms", "100"
};
assertEquals(100, new ConsoleConsumerOptions(validTimeoutMs).timeoutMs());
}
}
Loading