diff --git a/log4j-appender/src/main/java/org/apache/kafka/log4jappender/KafkaLog4jAppender.java b/log4j-appender/src/main/java/org/apache/kafka/log4jappender/KafkaLog4jAppender.java index 8ea0af337c12c..97a1ec96218fb 100644 --- a/log4j-appender/src/main/java/org/apache/kafka/log4jappender/KafkaLog4jAppender.java +++ b/log4j-appender/src/main/java/org/apache/kafka/log4jappender/KafkaLog4jAppender.java @@ -35,9 +35,11 @@ import static org.apache.kafka.clients.CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG; import static org.apache.kafka.clients.CommonClientConfigs.SECURITY_PROTOCOL_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.ACKS_CONFIG; +import static org.apache.kafka.clients.producer.ProducerConfig.BATCH_SIZE_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.COMPRESSION_TYPE_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG; +import static org.apache.kafka.clients.producer.ProducerConfig.LINGER_MS_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.MAX_BLOCK_MS_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.RETRIES_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG; @@ -76,6 +78,8 @@ public class KafkaLog4jAppender extends AppenderSkeleton { private int retries = Integer.MAX_VALUE; private int requiredNumAcks = 1; private int deliveryTimeoutMs = 120000; + private int lingerMs = 0; + private int batchSize = 16384; private boolean ignoreExceptions = true; private boolean syncSend; private Producer producer; @@ -100,6 +104,22 @@ public void setRequiredNumAcks(int requiredNumAcks) { this.requiredNumAcks = requiredNumAcks; } + public int getLingerMs() { + return lingerMs; + } + + public void setLingerMs(int lingerMs) { + this.lingerMs = lingerMs; + } + + public int getBatchSize() { + return batchSize; + } + + public void setBatchSize(int batchSize) { + this.batchSize = batchSize; + } + public int getRetries() { return retries; } @@ -269,6 +289,8 @@ public void activateOptions() { props.put(ACKS_CONFIG, Integer.toString(requiredNumAcks)); props.put(RETRIES_CONFIG, retries); props.put(DELIVERY_TIMEOUT_MS_CONFIG, deliveryTimeoutMs); + props.put(LINGER_MS_CONFIG, lingerMs); + props.put(BATCH_SIZE_CONFIG, batchSize); if (securityProtocol != null) { props.put(SECURITY_PROTOCOL_CONFIG, securityProtocol);