📨 Section 15 · Question #29

Can You Read Messages from a Specific Partition

In Apache Kafka, a consumer can consume messages directly from explicitly specified partitions of a topic. Instead of dynamic group subscription via consumer.subscribe(), you us...


🟢 Junior Level

Direct Answer: Yes, Absolutely!

In Apache Kafka, a consumer can consume messages directly from explicitly specified partitions of a topic. Instead of dynamic group subscription via consumer.subscribe(), you use manual partition assignment via consumer.assign().

Dynamic Group Subscription (subscribe):
Consumer ──subscribe("orders")──> Broker Group Coordinator decides partition allocation 
                                  (e.g., assigns Partitions 0 and 1 dynamically).

Manual Direct Assignment (assign):
Consumer ──assign([orders-0])───> Reads STRICTLY and ONLY Partition 0.
                                  Broker group coordination is completely bypassed.

Basic Java Example

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);

// 1. Explicitly specify target partition (Partition 0 of "orders")
TopicPartition partition0 = new TopicPartition("orders", 0);
consumer.assign(List.of(partition0));

// 2. Poll messages in standard event loop
while (running) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        System.out.printf("Partition: %d, Offset: %d, Key: %s, Value: %s%n",
            record.partition(), record.offset(), record.key(), record.value());
    }
}

Spring Kafka Declarative Partition Assignment

In Spring Boot applications, explicit partition targeting is configured declaratively using the @KafkaListener annotation:

@KafkaListener(
    topicPartitions = @org.springframework.kafka.annotation.TopicPartition(
        topic = "orders", 
        partitions = { "0", "2" } // Consumes exclusively from partitions 0 and 2
    )
)
public void listenSpecificPartitions(ConsumerRecord<String, String> record) {
    log.info("Received from Partition {}: {}", record.partition(), record.value());
}

🟡 Middle Level

Fundamental Differences: subscribe() vs assign()

The subscribe() and assign() methods are mutually exclusive. Attempting to invoke both methods on the same KafkaConsumer instance throws an IllegalStateException.

Characteristic consumer.subscribe() consumer.assign()
Group Coordinator Active participant via GroupCoordinator Completely bypassed
Rebalancing Automated failover and dynamic scaling Disabled. Consumer operates standalone
Fault Tolerance High: dead consumer’s partition is reassigned Zero: if consumer crashes, partition stalls
group.id Requirement Mandatory Optional (only required to commit offsets)
Dynamic Partitions Discovered automatically via metadata Ignored (requires manual code update)
Primary Use Cases Standard elastic microservice consumers Replays, audits, diagnostics, pinned state stores

Controlling Fetch Offsets via the Seek API

The true power of manual assignment is unlocked when paired with explicit positioning via Kafka’s Seek API:

TopicPartition tp = new TopicPartition("orders", 0);
consumer.assign(List.of(tp));

// 1. Seek to the absolute beginning of the partition (LogStartOffset)
consumer.seekToBeginning(List.of(tp));

// 2. Seek to the tail of the log (LogEndOffset) to process only incoming records
consumer.seekToEnd(List.of(tp));

// 3. Seek to a specific numerical offset
consumer.seek(tp, 450200L);

// 4. Time-based seek (seek to the first offset recorded at or after a specific timestamp)
long oneHourAgo = System.currentTimeMillis() - TimeUnit.HOURS.toMillis(1);
Map<TopicPartition, Long> timestampsToSearch = Map.of(tp, oneHourAgo);

Map<TopicPartition, OffsetAndTimestamp> foundOffsets = consumer.offsetsForTimes(timestampsToSearch);
OffsetAndTimestamp target = foundOffsets.get(tp);

if (target != null) {
    consumer.seek(tp, target.offset());
}

🔴 Senior Level

Offset Management and __consumer_offsets with assign()

A prevalent interview myth is: “When using assign(), group.id is useless and committing offsets to Kafka is unsupported.”

This is incorrect. The exact runtime behavior depends on consumer configuration:

  1. group.id is NOT configured:
    • Calling consumer.commitSync() or consumer.commitAsync() throws InvalidGroupIdException: To use the group management or offset commit APIs, you must provide a valid group.id in the consumer configuration.
    • The consumer does not persist state to the broker. On restart, it relies strictly on auto.offset.reset unless programmatically positioned via seek().
  2. group.id IS configured:
    • The consumer does not participate in consumer group rebalance protocols (no JoinGroup / SyncGroup network requests).
    • However, calling commitSync() commits offsets into the internal __consumer_offsets topic under the configured group.id!
    • Upon application restart, consumer.assign() reads the last committed offset from __consumer_offsets and resumes seamlessly without duplicates.

Production Pattern: Kubernetes StatefulSet + Partition Pinning

In high-throughput analytics engines and streaming architectures (Apache Flink, Kafka Streams, embedded RocksDB), dynamic rebalance pauses are unacceptable. Systems adopt Deterministic Partition Pinning:

Kubernetes StatefulSet (order-worker):
Pod 0 (HOSTNAME="order-worker-0") ──> assign(TopicPartition("orders", 0)) ──> Local Embedded RocksDB 0
Pod 1 (HOSTNAME="order-worker-1") ──> assign(TopicPartition("orders", 1)) ──> Local Embedded RocksDB 1
Pod 2 (HOSTNAME="order-worker-2") ──> assign(TopicPartition("orders", 2)) ──> Local Embedded RocksDB 2

Architectural Advantages

  1. Zero Rebalance Outages: Rolling updates and deployments generate zero consumer group rebalances and zero Stop-the-World processing pauses.
  2. Data Locality: Local disk state (embedded RocksDB) remains pinned to its designated partition without expensive state transfers across the cluster network.
  3. Deterministic Isolation: Crash logs, partition lag, and memory profiles are isolated cleanly per container.

The Tradeoff

If Pod 1 crashes due to an Out-Of-Memory (OOM) error, Kubernetes must restart Pod 1. While Pod 1 is rebooting, Partition 1 data is not consumed by any other worker, as standalone pods cannot adopt neighboring partitions automatically.


4 Tricky Questions

1. What happens if an application calls consumer.subscribe(List.of("topicA")) and subsequently calls consumer.assign(List.of(new TopicPartition("topicB", 0))) on the same consumer instance?

Answer: The invocation throws IllegalStateException: java.lang.IllegalStateException: Subscription to topics, partitions and pattern are mutually exclusive. Internally, KafkaConsumer tracks assignment state in SubscriptionState (NONE, AUTO_TOPICS, AUTO_PATTERN, or USER_ASSIGNED). Transitioning between dynamic subscription and manual assignment requires an explicit call to consumer.unsubscribe(), which resets the internal subscription state and detaches the consumer from the broker group coordinator.

2. A Kafka cluster administrator scales a topic from 4 partitions to 8 on a live production cluster. How do consumers using subscribe() vs assign() respond?

Answer:

  1. subscribe() consumers: The background metadata thread periodically updates cluster topology (metadata.max.age.ms, default: 5 minutes). Upon discovering new partitions, the group coordinator initiates a group rebalance, automatically distributing partitions 4 through 7 among existing consumer group members.
  2. assign() consumers: They will never consume from the new partitions. The partition assignment list is stored statically in client memory. The application will continue reading only partitions 0 through 3. To read the new partitions, the application must either restart or implement programmatic polling via consumer.partitionsFor("topic") to dynamically re-invoke consumer.assign(updatedPartitions).

3. Will ConsumerRebalanceListener callbacks fire when using consumer.assign()?

Answer: No, never. The ConsumerRebalanceListener interface (with onPartitionsRevoked() and onPartitionsAssigned()) is only accepted as a parameter in subscribe() overloads. The assign() method takes only a Collection<TopicPartition> and has no listener overloads. Because the consumer bypasses the broker group coordination protocol entirely, no partition revocation or assignment events are ever triggered.

4. How can you implement strict Exactly-Once processing with a relational database using assign() instead of subscribe()?

Answer: By using the Atomic Dual-Store Offsets pattern:

  1. The consumer binds to a specific partition via assign(tp).
  2. The consumer polls records and processes business logic within a local database transaction.
  3. In the very same database transaction, the application updates both the business table and an offset tracking table:
    BEGIN TRANSACTION;
    INSERT INTO orders (id, customer_id, amount) VALUES (4012, 99, 250.00);
    UPDATE kafka_offsets SET committed_offset = 105 WHERE partition_id = 0;
    COMMIT;
    
  4. On startup or after recovering from a crash:
    long lastCommitted = db.query("SELECT committed_offset FROM kafka_offsets WHERE partition_id = 0");
    consumer.seek(tp, lastCommitted + 1);
    
  5. Because the business entity and the consumed offset are committed atomically inside the relational database, processing is strictly Exactly-Once without distributed 2PC or Kafka Transactions overhead.

🎯 Interview Cheat Sheet

30-Second Elevator Pitch

“Yes, consumers can read from specific partitions using consumer.assign(Collection<TopicPartition>). Unlike subscribe(), assign() completely disables consumer group coordination and automated rebalancing—running the consumer in standalone mode. This grants deterministic control over partitions and enables precision navigation using the Seek API (seek(), seekToBeginning(), offsetsForTimes()). The tradeoff is zero automated failover: if an assigned consumer crashes, its partition stalls until the pod restarts. Common production use cases include offset replaying, diagnostic debugging, time-travel auditing, and stateful architectures with pinned local caches (Kubernetes StatefulSets + RocksDB).”

Key Production Monitoring Metrics

  • kafka.consumer:type=consumer-fetch-manager-metrics,client-id={id} -> records-lag: Monitored per assigned partition to ensure individual workers keep pace.
  • Container process health in Kubernetes: Essential alert trigger since no neighboring consumer can adopt partitions of a crashed standalone worker.

Red Flags (DO NOT Say)

  • ❌ “You can call subscribe() and assign() on the same consumer to combine features.” (They are mutually exclusive and throw IllegalStateException).
  • ❌ “When using assign(), committing offsets to Kafka is impossible.” (Supported if group.id is configured in consumer properties).
  • ❌ “If an assign() consumer dies, Kafka reassigns the partition to an available pod.” (Zero automated failover exists under manual assignment).
  • ❌ “Partitions added at runtime automatically appear in assign().” (The assignment list is static in memory and ignores runtime partition scaling).