Skip to content
JackSparrow414
Go back

Kafka (Part 3): Sending JSON with a Shared Serializer and Improving Producer Throughput

Table of contents

Open Table of contents

Producer Sending Strategy

How can we ensure that a correctly formatted message eventually reaches Kafka? The approach here is:

  1. Use the default ack setting, all, and configure min.insync.replicas to improve fault tolerance.
  2. Enable retries.
  3. If retries fail within a certain period, save the message to the database and let a scheduled job continue sending it later.
  4. Certain exceptional conditions will produce duplicate messages. How to ensure each is consumed only once will be discussed in the consumer implementation.
  5. Here, we only need to ensure that a message always goes to the same partition, regardless of retries. Kafka selects the partition using hash(key) if a key is provided. Each of our messages has an ID that never changes across retries, and we will not change the partition count during peak traffic. These conditions ensure that every attempt sends the message to the same partition.

Developers familiar with message queues probably know that a message sent to a broker is not necessarily persisted to disk immediately. It is generally written to the operating system cache, and the OS decides when to flush it to disk. Kafka works the same way. Strictly speaking, the five steps above therefore do not guarantee that data sent to Kafka cannot be lost. Kafka provides log.flush options, flush.message, and flush.ms to control disk flushing and log writes. See the Broker Configuration and Topic Configuration sections of the official documentation for details.

Kafka recommends leaving these settings alone and letting the operating system decide. The documentation states:

We generally feel that the guarantees provided by replication are stronger than sync to local disk, however the paranoid still may prefer having both and application level fsync policies are still supported

In other words, data reliability is ensured through replication rather than forced local disk flushes.

If you have only one standalone node with no replicas, however, you can consider disk-flush settings.

For more on flushing, read Application vs. OS Flush Management and the following sections in the Kafka documentation.

Using a Shared Serializer

Our message format is JSON, so Jackson must serialize classes to JSON strings. If we have several types of POJO messages, implementing the official Serializer interface separately for each is not ideal. Could a generic implementation do this for us? The open-source Kafka documentation does not cover this, but I found an implementation in the official Confluent GitHub repository. Simply add the dependency:

        <dependency>
            <groupId>io.confluent</groupId>
            <artifactId>kafka-json-serializer</artifactId>
            <version>7.5.1</version>
        </dependency>

Also specify the repository URL:

    <repositories>
        <repository>
            <id>confluent</id>
            <url>https://packages.confluent.io/maven/</url>
        </repository>
    </repositories>

The configuration is this line in the detailed code below:

result.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaJsonSerializer.class.getName());

Configuring Producer Options

A few points need attention:

  1. Enable compression.
  2. The documentation recommends not configuring retries directly. Use delivery.timeout.ms to control retry duration; its default is two minutes.
  3. Leave buffer.memory at its default of 32 MB unless there is a specific reason to change it.
  4. Use the default ack setting, all.
  5. Configure client.id to prevent InstanceAlreadyExistsException.
    /**
     * Tune these settings using the official documentation, relevant chapters of Kafka: The Definitive Guide, and your actual throughput requirements.
     * Locally, bootstrap.server IP addresses should match EXTERNAL in docker-compose.yml.
     * For message compression, the documentation recommends lz4; see https://www.confluent.io/blog/apache-kafka-message-compression/
     *
     * Set client.id to prevent InstanceAlreadyExistsException.
     * Otherwise Kafka generates a client.id, defaulting to producer-1; see {@link ProducerConfig#maybeOverrideClientId(Map)}.
     * Kafka Java Client uses client.id to generate a JMX ObjectName; see registerAppInfo in {@link KafkaProducer#KafkaProducer(ProducerConfig, Serializer, Serializer, ProducerMetadata, KafkaClient, ProducerInterceptors, Time)}.
     * If multiple applications (processes) omit client.id, IDs generated by the default rule can repeat, causing InstanceAlreadyExistsException.
     * Multiple producers in one application (one process) do not cause this when client.id is omitted, because of the incrementing counter {@link ProducerConfig#PRODUCER_CLIENT_ID_SEQUENCE}.
     */
    public static Properties loadProducerConfig(String valueSerializer) {
        Properties result = new Properties();
        result.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.0.102:9093");
        // Set client.id
        result.put(ProducerConfig.CLIENT_ID_CONFIG, SERVER_ID);
        result.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        result.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaJsonSerializer.class.getName());
        result.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, CompressionType.LZ4.name);
        // Each email message is about 20KB; default throughput is low, so these settings improve Kafka throughput
        // Default: 16384 bytes. Too small: email messages are sent one at a time rather than batched, which does not fit this scenario
        result.put(ProducerConfig.BATCH_SIZE_CONFIG, 1048576 * 10);
        // Default: 1048576 bytes. This limits batch size and is too small for messages of 20KB
        result.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 1048576 * 10);
        // Wait 10ms to accumulate more messages in a batch and improve throughput
        result.put(ProducerConfig.LINGER_MS_CONFIG, 10);
        return result;
    }

Improving Throughput

Sending Messages

@Log
public class MessageProducer {

    public static final KafkaProducer<String, UserDTO> PRODUCER = new KafkaProducer<>(KafkaConfiguration.loadProducerConfig(UserDTOSerializer.class.getName()));

    private MessageFailedService messageFailedService = new MessageFailedService();

    /**
     * Kafka producers retry failed sends. Related options are retries and delivery.timeout.ms; the documentation recommends delivery.timeout.ms, defaulting to two minutes.
     * The callback runs only after the final retry. For local retry testing, see https://lists.apache.org/thread/nwg326bxpo7ry116nqhxmsmc3bokc6hm
     * @param userDTO
     */
    public void sendMessage(final UserDTO userDTO) {
        ProducerRecord<String, UserDTO> user = new ProducerRecord<>("email", userDTO.getMessageId(),  userDTO);
        try {
            PRODUCER.send( user, (recordMetadata, e) -> {
                if (Objects.nonNull(e)) {
                    log.severe("message has sent failed");
                    MessageFailedEntity messageFailedEntity = new MessageFailedEntity();
                    messageFailedEntity.setMessageId(userDTO.getMessageId());
                    ObjectMapper mapper = new ObjectMapper();
                    try {
                        messageFailedEntity.setMessageContentJsonFormat(mapper.writeValueAsString(userDTO));
                    } catch (JsonProcessingException jsonProcessingException) {
                        log.severe("message content json format failed");
                    }
                    messageFailedEntity.setMessageType(MessageType.EMAIL);
                    messageFailedEntity.setMessageFailedPhrase(MessageFailedPhrase.PRODUCER);
                    messageFailedEntity.setFailedReason(e.getMessage());
                    // If sendMessage receives a list, the same applies: this must not be outside list.foreach
                    // If handled in the main thread, Kafka producer execution is asynchronous,
                    // so it may lag behind the main thread and return an empty value, e.g. an empty failedReason
                    messageFailedService.saveOrUpdateMessageFailed(messageFailedEntity);
                } else {
                    log.info("message has sent to topic: " + recordMetadata.topic() + ", partition: " + recordMetadata.partition() );
                }
            });
        } catch (TimeoutException e) {
            log.info("send message to kafka timeout, message: ");
            // TODO: Custom handling, e.g. email the Kafka administrator
        }
    }
}

A few explanations of the code:

  1. Messages are sent asynchronously. Log a message when sending succeeds.
  2. The key point is retry behavior: the callback runs only after the final retry, not once per retry. I asked the community by email; see this thread. To test or debug retries locally, increase min.insync.replicas. For example, with one Kafka node, set it above 1 and configure producer acknowledgments as all. Sending a message then produces retry logs.

Closing the Producer

Implement ServletContextListener and configure it in a listener element in web.xml:

public class KafkaListener implements ServletContextListener {

    private static final List<KafkaProducer> KAFKA_PRODUCERS = new LinkedList<>();

    @Override
    public void contextInitialized(ServletContextEvent sce) {
        KAFKA_PRODUCERS.add(MessageProducer.PRODUCER);
    }

    @Override
    public void contextDestroyed(ServletContextEvent sce) {
        KAFKA_PRODUCERS.forEach(KafkaProducer::close);
    }
}
<?xml version="1.0" encoding="UTF-8" ?>
<web-app xmlns="https://jakarta.ee/xml/ns/jakartaee"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="https://jakarta.ee/xml/ns/jakartaee
                      https://jakarta.ee/xml/ns/jakartaee/web-app_6_0.xsd"
         version="6.0">

  <listener>
    <listener-class>com.business.server.listener.KafkaListener</listener-class>
  </listener>
</web-app>

Closing Thoughts

  1. While implementing this, consult the relevant chapters of Kafka: The Definitive Guide, or cloud providers’ Kafka developer documentation. I recommend the book. Although Alibaba Cloud and Huawei Cloud claim compatibility with open-source Kafka, I found their versions lag behind, and many best practices are outdated.
  2. There is nothing particularly unusual on the producer side. The main tasks are designing the message format for the business scenario and minimizing message size.
  3. If your messages are larger than mine, exceeding 1 MB, both producer throughput and consumer processing speed become issues. Without a concrete scenario, I cannot suggest a good approach. The option I can think of is reducing serialized message size, perhaps with Avro or Protobuf, though I have not tried either. Please share your experience if you have used them.

Sample Source Repository

  1. GitHub repository
  2. The business-server module represents the producer.
  3. IDEA run configuration:Deployment tab of the Tomcat 10.1.11 run configuration in IDEA: deploys business-server:war exploded with Application context set to /business-server Pay attention to Application context. After startup, visit the port plus Application context, for example:
http://localhost:8999/business-server

The next post covers consuming messages, important consumer configuration options, and retry mechanisms in the consumption logic.


Share this post:

Continue this series

Kafka in Practice

  1. Kafka (Part 1): A Single-Node KRaft Setup with Docker Compose, Kafka UI, and Prometheus JMX Exporter
  2. Kafka (Part 2): Designing a Messaging System to Decouple Email Delivery
  3. Kafka (Part 3): Sending JSON with a Shared Serializer and Improving Producer ThroughputYou are here
  4. Kafka (Part 4): Consuming JSON, Sharing a Deserializer, and Improving Throughput
  5. Kafka (Part 5): Consumer Callbacks, Scheduled Retries, and Rebalancing
  6. Kafka (Part 6): Oracle-to-PostgreSQL CDC with Kafka Connect and Debezium, and Cache Consistency
  7. Kafka (Part 7): Integrating Apache Avro and Apicurio Schema Registry to Ensure Message Compatibility Between Producers and Consumers

Comments

Questions, corrections, and experiences are welcome. Sign in with GitHub to comment; both language versions share this discussion.

Comments are available on the live site only.