Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

A basic Kafka consumer in Java needs a broker address, a consumer group ID, key and value deserializers, a topic subscription, and a loop that calls poll(). The key design choice is when to commit offsets: automatic commits are simpler, while committing after processing gives the application tighter control over which records it acknowledges.

Build a basic Java Kafka consumer

This example uses the Apache Kafka Java client, string keys and values, a consumer group named orders-consumer, and a topic named orders. It disables automatic commits and commits after handling each batch returned by poll().

import java.time.Duration;
import java.util.List;
import java.util.Properties;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

public class OrdersConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "orders-consumer");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
                  StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
                  StringDeserializer.class.getName());
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");

        try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(List.of("orders"));

            while (true) {
                ConsumerRecords<String, String> records =
                    consumer.poll(Duration.ofMillis(1000));

                for (ConsumerRecord<String, String> record : records) {
                    process(record.key(), record.value());
                }

                consumer.commitSync();
            }
        }
    }

    private static void process(String key, String value) {
        System.out.println("key=" + key + ", value=" + value);
    }
}

The example follows the Apache Kafka trunk consumer example. Check the client version used by your project before copying APIs or relying on defaults; the cited API documentation is for Kafka 2.8.1, and its configuration reference is for Kafka 2.6. KafkaConsumer API, Kafka 2.8.1 · Kafka 2.6 consumer configuration · Apache Kafka trunk consumer example.

Configure the consumer before subscribing

Broker and group

bootstrap.servers is the initial broker address the client uses to connect to the Kafka cluster. Replace localhost:9092 with an address reachable from the application. group.id identifies the consumer group. Consumers using the same group ID share topic partitions, which enables parallel work across members and allows the group to redistribute partitions when membership changes.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Deserializers and reset behavior

Kafka stores record keys and values as bytes. The configured deserializers turn those bytes into the Java types used in KafkaConsumer<K,V>. The example uses StringDeserializer for both; use deserializers that match the actual record format if keys or values are not strings.

auto.offset.reset determines where a group begins only when it has no valid committed offset. The example sets it to earliest, so the consumer starts from the earliest available offset in that case. It does not rewind a group that already has a committed offset.

Subscribe, poll, and process records

subscribe() asks Kafka’s group management to assign partitions for the named topics. The application then repeatedly calls poll(Duration); each call returns a ConsumerRecords batch, which can contain zero or more records. The example polls with a one-second duration and processes every returned record before committing.

With group management, the application must continue calling poll() within max.poll.interval.ms, the maximum delay between poll invocations. If processing takes too long and that interval is exceeded, Kafka can consider the consumer no longer active and trigger a rebalance. max.poll.records limits how many records one poll returns. The Kafka 2.6 configuration reference lists defaults of 300000 milliseconds for max.poll.interval.ms and 500 for max.poll.records; verify values for the client version in use rather than treating those as universal defaults. Kafka 2.6 consumer configuration

What’s actually slowing this PC down?

Pick the symptom - the matching free tool is one click away.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Choose when offsets are committed

Approach When the offset is committed Practical effect
Automatic commit Periodically in the background when enable.auto.commit=true. Simpler to configure, but the committed position is not necessarily tied to successful completion of your processing logic.
Manual synchronous commit When the application calls commitSync(); in the example, after processing a poll batch. Gives the application control over commit timing. The call blocks and reports unrecoverable errors.
Manual asynchronous commit When the application calls commitAsync(). Does not block; commit errors are reported through a callback, so the application must decide how to handle them.

A committed offset is the next record the group should consume, not the last record already handled. For a successfully processed record at offset n, the corresponding committed position is n + 1. The Kafka API documentation expresses this as “lastProcessedMessageOffset + 1.” KafkaConsumer API, offset management

When processing correctness requires tighter control, set enable.auto.commit=false and commit only after successful processing. If processing fails before the commit, the record may be delivered again after restart or reassignment, so make processing idempotent or otherwise safe to retry. Committing before processing can instead leave the group past work that did not complete.

Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

Choose group-managed subscription or explicit assignment

Use subscribe() when you want Kafka’s consumer group management to assign partitions and rebalance them as group membership changes. This is the normal choice for a group of cooperating application instances.

Explicit assignment with assign() lets an application name the partitions it will read rather than relying on group assignment. It is useful when the application must control partition ownership itself, but it also makes partition assignment and reassignment the application’s responsibility. Do not mix subscribe() and assign() in the same consumer’s assignment strategy.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Choose which transactional records are visible

By default, isolation.level is read_uncommitted, which allows the consumer to see records from transactions that may later be aborted. Setting it to read_committed hides aborted transactional records. Choose the setting to match whether the application must observe only committed transactional data. Kafka 2.6 consumer configuration

Prepare the teaching example for an application

The loop demonstrates the core API flow, but a deployed consumer needs application-specific handling around it:

  • Shutdown: arrange a controlled way to stop the polling loop and close the consumer, rather than relying on an unconditional infinite loop.
  • Failures: define what happens when record processing or a commit fails, including whether to retry, stop, or route a record to a dead-letter workflow.
  • Retries and duplicates: because an uncommitted record can be consumed again, make side effects idempotent where possible or design another application-level safeguard.
  • Processing duration: ensure the time spent processing batches remains compatible with max.poll.interval.ms; if work can take longer, design the processing and polling approach accordingly.

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.