From d65e8ff59dadcbe4ede5d14907a1ecf708cb1705 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 16 Jan 2019 12:47:14 -0500 Subject: [PATCH] GH-529: Add lz4 and docs for zstd compression Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/529 - Also log exception when getting partition information --- docs/src/main/asciidoc/overview.adoc | 7 +++++++ .../properties/KafkaProducerProperties.java | 17 ++++++++++++++++- .../provisioning/KafkaTopicProvisioner.java | 1 + 3 files changed, 24 insertions(+), 1 deletion(-) diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index 667ddc393..816047f37 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -317,6 +317,13 @@ If a topic already exists with a smaller partition count and `autoAddPartitions` If a topic already exists with a smaller partition count and `autoAddPartitions` is enabled, new partitions are added. If a topic already exists with a larger number of partitions than the maximum of (`minPartitionCount` or `partitionCount`), the existing partition count is used. +compression:: +Set the `compression.type` producer property. +Supported values are `none`, `gzip`, `snappy` and `lz4`. +If you override the `kafka-clients` jar to 2.1.0 (or later), as discussed in the https://docs.spring.io/spring-kafka/docs/2.2.x/reference/html/deps-for-21x.html[Spring for Apache Kafka documentation], and wish to use `zstd` compression, use `spring.cloud.stream.kafka.bindings..producer.configuration.compression.type=zstd`. ++ +Default: `none`. + ==== Usage examples In this section, we show the use of the preceding properties for specific scenarios. diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java index b722bc89c..c4444b981 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java @@ -117,17 +117,32 @@ public class KafkaProducerProperties { * Enumeration for compression types. */ public enum CompressionType { + /** * No compression. */ none, + /** * gzip based compression. */ gzip, + /** * snappy based compression. */ - snappy + snappy, + + /** + * lz4 compression + */ + lz4, + + // /** // TODO: uncomment and fix docs when kafka-clients 2.1.0 or newer is the default + // * zstd compression + // */ + // zstd + } + } diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index 4749340e0..e39a6294f 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -414,6 +414,7 @@ public class KafkaTopicProvisioner implements ProvisioningProvider