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:
| Option | What 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
--propertyoption is available for both the console producer and consumer, but this should NOT be confused with--producer-propertyor--consumer-property. THE--propertyOPTION IS ONLY USED FOR PASSING CONFIGURATIONS TO THE MESSAGE FORMATTER, AND NOT THE CLIENT ITSELF."
Useful producer options:
| Option | Meaning |
|---|---|
--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):
| Property | Meaning |
|---|---|
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
LineMessageReaderwill split the input on the FIRST instance of thekey.separator. If there are no characters remaining after that, the value of the message will be EMPTY. If no key separator is present, or ifparse.keyisfalse, the key will benull."
CHANGING LINE-READING BEHAVIOR
"You can provide your own class... must extend
kafka.common.MessageReaderand will be responsible for creating theProducerRecord. 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:
| Option | What it selects |
|---|---|
--topic | a 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:
| Option | Meaning |
|---|---|
--formatter <classname> | Message formatter class. Default kafka.tools.DefaultMessageFormatter |
--from-beginning | Consume from the oldest offset; otherwise from the latest |
--max-messages <int> | Max messages before exiting |
--partition <int> | Consume only from that partition |
--offset | An 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:
| Formatter | What 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.partitionkey.separator·line.separatorkey.deserializer·value.deserializer
"The deserializer classes must implement
org.apache.kafka.common.serialization.Deserializer, and the console consumer will calltoStringon them to get the output. Typically you would insert them into the classpath by setting theCLASSPATHenvironment 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.)