Learn Labs
12. Administering Kafka

12.4 Producing and consuming from the console

Configs via --consumer.config <file> or --consumer-property <key>=<value>.

⚠️ PIPING OUTPUT TO ANOTHER APPLICATION

"While it is possible to write applications that wrap around the console consumer or producer... this type of application is QUITE FRAGILE and SHOULD BE AVOIDED. IT IS DIFFICULT TO INTERACT WITH THE CONSOLE CONSUMER IN A WAY THAT DOES NOT LOSE MESSAGES. Likewise, the console producer DOES NOT ALLOW FOR USING ALL FEATURES, and properly sending bytes is tricky. IT IS BEST TO USE EITHER THE JAVA CLIENT LIBRARIES DIRECTLY or a third-party client library."

4.1 Console producer

kafka-console-producer.sh --bootstrap-server localhost:9092 --topic my-topic
>Message 1
>Test Message 2
>^D          ← EOF (Ctrl-D) closes the client

"By default, messages are read one per line, with a TAB CHARACTER separating the key and the value (if no tab character is present, the key is null). ... the producer reads in and produces raw bytes using the default serializer (DefaultEncoder)."

Passing producer configs — two ways:

OptionWhat you give it
--producer.config <config-file>a full properties file
--producer-property <key>=<value>one or more on the command line

► “useful for producer options like message-batching configurations (such as linger.ms or batch.size)”

⚠️ CONFUSING COMMAND-LINE OPTIONS

"The --property option is available for both the console producer and consumer, but this should NOT be confused with --producer-property or --consumer-property. THE --property OPTION IS ONLY USED FOR PASSING CONFIGURATIONS TO THE MESSAGE FORMATTER, AND NOT THE CLIENT ITSELF."

Useful producer options:

OptionMeaning
--batch-size"number of messages sent in a single batch if they are not being sent synchronously"
--timeout"If a producer is running in asynchronous mode, max time waiting for the batch size before producing — to avoid long waits on low-producing topics"
--compression-codec <string>none, gzip, snappy, zstd, or lz4. Default: gzip
--sync"Produce messages synchronously, waiting for each message to be acknowledged before sending the next"

Line-reader options (via --property, for kafka.tools.ConsoleProducer$LineMessageReader):

PropertyMeaning
ignore.error"Set to false to throw an exception when parse.key is true and a key separator is not present. Defaults to true."
parse.key"Set to false to always set the key to null. Defaults to true."
key.separator"delimiter between key and value. Defaults to a tab character."

"the LineMessageReader will split the input on the FIRST instance of the key.separator. If there are no characters remaining after that, the value of the message will be EMPTY. If no key separator is present, or if parse.key is false, the key will be null."

CHANGING LINE-READING BEHAVIOR

"You can provide your own class... must extend kafka.common.MessageReader and will be responsible for creating the ProducerRecord. Specify it with --line-reader, and make sure the JAR containing your class is in the classpath."

4.2 Console consumer

kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --whitelist 'my.*' --from-beginning

"messages are printed in standard output, delimited by a new line. By default, it outputs the raw bytes in the message, WITHOUT THE KEY, with no formatting (using DefaultFormatter)."

Two mutually exclusive topic selectors:

OptionWhat it selects
--topica Single topic
--whitelist“A regular expression matching all topics to consume from (Remember to properly escape the regex so that it is not processed improperly by the shell)”

► “Only One of the previous options should be selected and used.”

Useful consumer options:

OptionMeaning
--formatter <classname>Message formatter class. Default kafka.tools.DefaultMessageFormatter
--from-beginningConsume from the oldest offset; otherwise from the latest
--max-messages <int>Max messages before exiting
--partition <int>Consume only from that partition
--offsetAn offset <int>, or earliest, or latest
--skip-message-on-error"Skip a message if there is an error when processing instead of halting. Useful for debugging."

Configs via --consumer.config <file> or --consumer-property <key>=<value>.

The four message formatters:

FormatterWhat it prints
kafka.tools.DefaultMessageFormatter(default)
kafka.tools.LoggingMessageFormatter“Outputs messages using the Logger, rather than standard out. Printed at the Info level and include the Timestamp, Key, and Value.”
kafka.tools.ChecksumMessageFormatter“Prints Only message checksums.”
kafka.tools.NoOpMessageFormatter“Consumes messages but Does not output them at all.”

DefaultMessageFormatter properties (Table 12-4), via --property:

  • print.timestamp · print.key · print.offset · print.partition
  • key.separator · line.separator
  • key.deserializer · value.deserializer

"The deserializer classes must implement org.apache.kafka.common.serialization.Deserializer, and the console consumer will call toString on them to get the output. Typically you would insert them into the classpath by setting the CLASSPATH environment variable before executing the script."

💡 4.3 Consuming __consumer_offsets — a real diagnostic

Why: "You may want to see if a particular group is committing offsets AT ALL, or HOW OFTEN offsets are being committed."

kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic __consumer_offsets --from-beginning --max-messages 1 \
  --formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter" \
  --consumer-property exclude.internal.topics=false

[my-group-name,my-topic,0]::[OffsetMetadata[1,NO_METADATA]
  CommitTime 1623034799990 ExpirationTime 1623639599990]

Three things this command requires that are easy to miss: the special formatter class (with the $ escaped), exclude.internal.topics=false, and --max-messages so you don't drown.

(Note the ExpirationTime in the output — that's offsets.retention.minutes from Ch. 4 §6.8 made visible.)


On this page