Kafka clients in production are not just "send a message" and "receive a message". Often, incorrect configuration compromises business data consistency. We recently encountered a project where auto-commit and lack of idempotency caused up to 10% of orders to be lost during consumer restarts. After implementing manual offset management and an idempotent producer, losses were reduced to zero. The cost of such losses for an average online store can reach 2 million rubles per month. We specialize in developing reliable producers and consumers for high-load web applications in Java and Python. This article covers key problems, typical configurations, and practical cases—from idempotency setup to parallel processing with order preservation.
Problems We Solve
Message Loss During Broker or Consumer Failures
Without proper configuration of acks and idempotency, you risk losing critical business data. For example, with acks=1, a message can be lost if the leader fails before replication. We use acks=all with min.insync.replicas=2 and enable.idempotence=true. This eliminates loss even if one broker fails. Infrastructure savings by avoiding expensive solutions: up to 40%.
Order Violation and Duplication
Consumers without offset management may process the same message twice, and parallel processing destroys order. The solution: manual commit after processing and routing by key within a partition queue. On one project, this reduced validation errors by 30%.
Long Rebalance and Re-processing
When a consumer fails and a rebalance occurs, the group may re-read large amounts of data. We configure session.timeout.ms and max.poll.interval.ms so rebalance takes no more than 5 seconds, and we use DLQ to avoid blocking on "bad" messages.
How We Guarantee Delivery
According to Apache Kafka documentation, delivery guarantees are determined by the combination of acks, idempotency, and transactions. Let's look at key patterns.
Idempotent Producer with acks=all
For financial transactions, order events, and any critical stream, we set enable.idempotence=true and acks=all. This eliminates duplicates even with retries. Throughput impact is minor (5-10%), but reliability is significantly improved.
// Java — idempotent producer
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);
props.put("schema.registry.url", "http://schema-registry:8081");
// Reliability
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);
props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120_000);
// Performance
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 65536); // 64KB batch
props.put(ProducerConfig.LINGER_MS_CONFIG, 5); // Wait up to 5ms to fill batch
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 67_108_864); // 64MB buffer
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
KafkaProducer<String, OrderEvent> producer = new KafkaProducer<>(props);
Sending with error handling:
public CompletableFuture<RecordMetadata> sendOrderEvent(OrderEvent event) {
ProducerRecord<String, OrderEvent> record = new ProducerRecord<>(
"order-events",
event.getOrderId(),
event
);
CompletableFuture<RecordMetadata> future = new CompletableFuture<>();
producer.send(record, (metadata, exception) -> {
if (exception != null) {
if (exception instanceof RetriableException) {
log.error("Retriable error, Kafka will retry: {}", exception.getMessage());
} else {
log.error("Fatal producer error for order {}: {}", event.getOrderId(), exception.getMessage());
future.completeExceptionally(exception);
}
} else {
log.debug("Sent to partition {} offset {}", metadata.partition(), metadata.offset());
future.complete(metadata);
}
});
return future;
}
Example consumer configuration with manual commit
Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092,kafka-2:9092,kafka-3:9092");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "order-processor");
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
consumerProps.put("schema.registry.url", "http://schema-registry:8081");
consumerProps.put("specific.avro.reader", true);
// Disable auto-commit
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
// Session timeout
consumerProps.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30_000);
consumerProps.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10_000);
consumerProps.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);
consumerProps.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300_000);
KafkaConsumer<String, OrderEvent> consumer = new KafkaConsumer<>(consumerProps);
consumer.subscribe(List.of("order-events"), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
commitCurrentOffsets();
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
log.info("Assigned partitions: {}", partitions);
}
});
Map<TopicPartition, OffsetAndMetadata> pendingOffsets = new HashMap<>();
try {
while (!shutdown.get()) {
ConsumerRecords<String, OrderEvent> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, OrderEvent> record : records) {
try {
processOrder(record.value());
pendingOffsets.put(
new TopicPartition(record.topic(), record.partition()),
new OffsetAndMetadata(record.offset() + 1)
);
} catch (NonRetriableException e) {
sendToDlq(record, e);
pendingOffsets.put(
new TopicPartition(record.topic(), record.partition()),
new OffsetAndMetadata(record.offset() + 1)
);
}
}
if (!pendingOffsets.isEmpty()) {
consumer.commitSync(pendingOffsets);
pendingOffsets.clear();
}
}
} finally {
consumer.close();
}
Why Manual Offset Management is Better Than Auto-Commit
auto-commit creates an illusion of simplicity but hides a risk: the offset is committed before processing. If the consumer crashes between commit and processing, the message is lost irretrievably. We always disable enable.auto.commit and commit manually after successful processing of the record. This is a standard practice for production systems.
Comparison of acks Modes
| Mode | acks | Data Loss on Leader Failure | Performance | Recommendation |
|---|---|---|---|---|
| Fire-and-forget | 0 | Possible | Maximum | Only for non-critical logs |
| Leader acknowledges | 1 | Possible (non-replicated) | High | Internal services with retries |
| All ISR | all (with min.insync.replicas) | Impossible | 10-15% lower | Critical data, orders, payments |
Comparison of Java and Python Clients
| Criterion | Java (kafka-clients) | Python (confluent-kafka) |
|---|---|---|
| Performance | High (native JVM) | Medium (via C library) |
| Transaction support | Full | Limited |
| Offset management | Manual + async | Mostly sync commit |
| Memory consumption | Higher (JVM heap) | Lower |
| Recommendation | Critical, high-throughput | Quick prototypes, medium load |
How to Organize Parallel Processing Without Losing Order
One poll thread + a pool of workers keyed by partition key is the classic pattern. We preserve order within a single order's events while parallelizing across orders.
Map<Integer, BlockingQueue<ConsumerRecord<String, OrderEvent>>> partitionQueues = new HashMap<>();
ExecutorService workers = Executors.newFixedThreadPool(12);
for (ConsumerRecord<String, OrderEvent> record : records) {
int partitionIndex = record.partition() % NUM_WORKERS;
workerQueues.get(partitionIndex).offer(record);
}
Turnkey Development Process
- Analytics — define topics, message keys, format (Avro/JSON/Protobuf), consumer groups.
- Implementation — write producers with idempotency, consumers with manual commit, DLQ.
- Testing — unit tests with EmbeddedKafka, integration tests with Testcontainers, load testing.
- Deployment — configure monitoring (consumer lag, errors) and alerts.
What's Included in the Work
- Repository with code (Java/Python) and configuration documentation.
- Docker images configured for staging/production.
- Deployment and rollback instructions.
- Training for the customer's team (1-2 hours).
- Code warranty: 3 months of free support.
Timeline
From 6 to 10 working days depending on complexity. Request a project evaluation — we'll respond within one day.
Common Mistakes (Checklist)
- [ ]
enable.idempotencenot enabled — duplicates on retries. - [ ] auto-commit used — message loss on consumer failure.
- [ ]
max.poll.recordstoo large — long rebalance. - [ ] No handling of
NonRetriableException— message lost without DLQ. - [ ] Mismatched serializer versions between producer and consumer — deserialization errors.
If you recognize your project in these issues, contact us. We have developed producers and consumers for 15+ projects handling up to 100,000 messages per second.
Guarantee: We fix deadlines in the contract and provide 3 months of post-delivery support.







