Custom Kafka Producers and Consumers Development

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 i

Development and maintenance of all types of websites:

Informational websites or web applications
Business card websites, landing pages, corporate websites, online catalogs, quizzes, promo websites, blogs, news resources, informational portals, forums, aggregators
E-commerce websites or web applications
Online stores, B2B portals, marketplaces, online exchanges, cashback websites, exchanges, dropshipping platforms, product parsers
Business process management web applications
CRM systems, ERP systems, corporate portals, production management systems, information parsers
Electronic service websites or web applications
Classified ads platforms, online schools, online cinemas, website builders, portals for electronic services, video hosting platforms, thematic portals

These are just some of the technical types of websites we work with, and each of them can have its own specific features and functionality, as well as be customized to meet the specific needs and goals of the client.

Our competencies:

Frequently Asked Questions

Latest works

  • image_website-b2b-advance_0.webp
    B2B ADVANCE company website development
    1419
  • image_web-applications_feedme_466_0.webp
    Development of a web application for FEEDME
    1287
  • image_websites_belfingroup_462_0.webp
    Website development for BELFINGROUP
    983
  • image_ecommerce_furnoro_435_0.webp
    Development of an online store for the company FURNORO
    1244
  • image_crm_enviok_479_0.webp
    Development of a web application for Enviok
    983
  • image_bitrix-bitrix-24-1c_fixper_448_0.webp
    Website development for FIXPER company
    998

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

  1. Analytics — define topics, message keys, format (Avro/JSON/Protobuf), consumer groups.
  2. Implementation — write producers with idempotency, consumers with manual commit, DLQ.
  3. Testing — unit tests with EmbeddedKafka, integration tests with Testcontainers, load testing.
  4. 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.idempotence not enabled — duplicates on retries.
  • [ ] auto-commit used — message loss on consumer failure.
  • [ ] max.poll.records too 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.