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:
group.idis NOT configured:- Calling
consumer.commitSync()orconsumer.commitAsync()throwsInvalidGroupIdException: 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.resetunless programmatically positioned viaseek().
- Calling
group.idIS configured:- The consumer does not participate in consumer group rebalance protocols (no
JoinGroup/SyncGroupnetwork requests). - However, calling
commitSync()commits offsets into the internal__consumer_offsetstopic under the configuredgroup.id! - Upon application restart,
consumer.assign()reads the last committed offset from__consumer_offsetsand resumes seamlessly without duplicates.
- The consumer does not participate in consumer group rebalance protocols (no
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
- Zero Rebalance Outages: Rolling updates and deployments generate zero consumer group rebalances and zero Stop-the-World processing pauses.
- Data Locality: Local disk state (embedded RocksDB) remains pinned to its designated partition without expensive state transfers across the cluster network.
- 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:
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.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 viaconsumer.partitionsFor("topic")to dynamically re-invokeconsumer.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:
- The consumer binds to a specific partition via
assign(tp). - The consumer polls records and processes business logic within a local database transaction.
- 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; - 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); - 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>). Unlikesubscribe(),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.idis 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).
Related Topics
- What is Consumer Group — Dynamic partition balancing
- How Does Consumer Balancing Work in a Group — Rebalance protocols
- What is Offset in Kafka — Offset manipulation and commits
- What is Rebalancing and When Does It Happen — Rebalance triggers