Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,18 @@ public interface ProducerBuilder {
*/
ProducerBuilder setTransactionChecker(TransactionChecker checker);

/**
* Set the threshold in bytes to enable message body compression.
*
* <p>Compression is disabled by default. Once this threshold is set, a message whose body size reaches it is
* compressed with GZIP before sending, and the consumer decompresses the body transparently according to the
* body encoding, so no change is required on the consuming side.
*
* @param compressBodyThresholdBytes threshold in bytes to enable message body compression.
* @return the producer builder instance.
*/
ProducerBuilder setCompressBodyThresholdBytes(int compressBodyThresholdBytes);

/**
* Finalize the build of {@link Producer} instance and start.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ public class ProducerBuilderImpl implements ProducerBuilder {
private final Set<String> topics = new HashSet<>();
private int maxAttempts = 3;
private TransactionChecker checker = null;
private int compressBodyThresholdBytes = Integer.MAX_VALUE;

public ProducerBuilderImpl() {
}
Expand Down Expand Up @@ -83,13 +84,24 @@ public ProducerBuilder setTransactionChecker(TransactionChecker checker) {
return this;
}

/**
* @see ProducerBuilder#setCompressBodyThresholdBytes(int)
*/
@Override
public ProducerBuilder setCompressBodyThresholdBytes(int compressBodyThresholdBytes) {
checkArgument(compressBodyThresholdBytes > 0, "compressBodyThresholdBytes should be positive");
this.compressBodyThresholdBytes = compressBodyThresholdBytes;
return this;
}

/**
* @see ProducerBuilder#build()
*/
@Override
public Producer build() {
checkNotNull(clientConfiguration, "clientConfiguration has not been set yet");
final ProducerImpl producer = new ProducerImpl(clientConfiguration, topics, maxAttempts, checker);
producer.publishingSettings.setCompressBodyThresholdBytes(compressBodyThresholdBytes);
producer.startAsync().awaitRunning();
return producer;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,10 @@ public class PublishingSettings extends Settings {
*/
private volatile int maxBodySizeBytes = 4 * 1024 * 1024;
private volatile boolean validateMessageType = true;
/**
* If message body size reaches this threshold, it is compressed with GZIP before sending. Disabled by default.
*/
private volatile int compressBodyThresholdBytes = Integer.MAX_VALUE;

public PublishingSettings(String namespace, ClientId clientId, Endpoints accessPoint,
ExponentialBackoffRetryPolicy retryPolicy, Duration requestTimeout, Set<String> topics,
Expand All @@ -57,6 +61,15 @@ public int getMaxBodySizeBytes() {
return maxBodySizeBytes;
}

public int getCompressBodyThresholdBytes() {
return compressBodyThresholdBytes;
}

// Client-only setting: configured via ProducerBuilder and never synced from the server, see sync().
void setCompressBodyThresholdBytes(int compressBodyThresholdBytes) {
this.compressBodyThresholdBytes = compressBodyThresholdBytes;
}

public boolean isValidateMessageType() {
return validateMessageType;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,13 +36,28 @@
public class PublishingMessageImpl extends MessageImpl {
private final MessageId messageId;
private final MessageType messageType;
private final Encoding encoding;
private final byte[] transportBody;

public PublishingMessageImpl(Message message, PublishingSettings publishingSettings, boolean txEnabled)
throws IOException {
super(message);
final int length = message.getBody().remaining();
// Compress the message body if it reaches the compression threshold. Fall back to the
// original body when compression does not reduce the size (e.g. already-compressed or
// encrypted payloads), which avoids inflating the transport body.
byte[] candidateBody = body;
Encoding candidateEncoding = Encoding.IDENTITY;
if (body.length >= publishingSettings.getCompressBodyThresholdBytes()) {
final byte[] compressedBody = Utilities.compressBytesGZIP(body);
if (compressedBody.length < body.length) {
candidateBody = compressedBody;
candidateEncoding = Encoding.GZIP;
}
}
this.encoding = candidateEncoding;
this.transportBody = candidateBody;
final int maxBodySizeBytes = publishingSettings.getMaxBodySizeBytes();
if (length > maxBodySizeBytes) {
if (transportBody.length > maxBodySizeBytes) {
throw new IOException("Message body size exceeds the threshold, max size=" + maxBodySizeBytes + " bytes");
}
// Generate message id.
Expand Down Expand Up @@ -114,7 +129,7 @@ public apache.rocketmq.v2.Message toProtobuf(String namespace, MessageQueueImpl
// Born host
.setBornHost(Utilities.hostName())
// Body encoding
.setBodyEncoding(Encoding.toProtobuf(Encoding.IDENTITY))
.setBodyEncoding(Encoding.toProtobuf(encoding))
// Queue id
.setQueueId(mq.getQueueId())
// Message type
Expand All @@ -136,7 +151,7 @@ public apache.rocketmq.v2.Message toProtobuf(String namespace, MessageQueueImpl
// Topic
.setTopic(topicResource)
// Message body
.setBody(ByteString.copyFrom(getBody()))
.setBody(ByteString.copyFrom(transportBody))
// System properties
.setSystemProperties(systemProperties)
// User properties
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,18 @@ public void testSetTransactionChecker() {
builder.setTransactionChecker(messageView -> TransactionResolution.COMMIT);
}

@Test(expected = IllegalArgumentException.class)
public void testSetNegativeCompressBodyThresholdBytes() {
final ProducerBuilderImpl builder = new ProducerBuilderImpl();
builder.setCompressBodyThresholdBytes(-1);
}

@Test
public void testSetCompressBodyThresholdBytes() {
final ProducerBuilderImpl builder = new ProducerBuilderImpl();
builder.setCompressBodyThresholdBytes(4096);
}

@Test(expected = NullPointerException.class)
public void testBuildWithoutClientConfiguration() {
final ProducerBuilderImpl builder = new ProducerBuilderImpl();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.rocketmq.client.java.impl.producer;

/**
* Test-only accessor for package-private mutators of {@link PublishingSettings}.
*/
public final class PublishingSettingsTestHelper {
private PublishingSettingsTestHelper() {
}

public static void setCompressBodyThresholdBytes(PublishingSettings settings, int compressBodyThresholdBytes) {
settings.setCompressBodyThresholdBytes(compressBodyThresholdBytes);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.rocketmq.client.java.message;

import static org.junit.Assert.assertArrayEquals;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;

import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
import java.util.Random;
import org.apache.rocketmq.client.apis.message.Message;
import org.apache.rocketmq.client.java.impl.producer.PublishingSettings;
import org.apache.rocketmq.client.java.impl.producer.PublishingSettingsTestHelper;
import org.apache.rocketmq.client.java.misc.Utilities;
import org.apache.rocketmq.client.java.route.MessageQueueImpl;
import org.apache.rocketmq.client.java.tool.TestBase;
import org.junit.Test;

public class PublishingMessageImplTest extends TestBase {

@Test
public void testBodyNotCompressedByDefault() throws IOException {
String topic = "testTopic";
byte[] body = "foobar".getBytes(StandardCharsets.UTF_8);
final Message message = new MessageBuilderImpl().setTopic(topic).setBody(body).build();
final PublishingSettings settings = fakeProducerSettings();
final PublishingMessageImpl publishingMessage = new PublishingMessageImpl(message, settings, false);
final MessageQueueImpl mq = fakeMessageQueueImpl(topic);
final apache.rocketmq.v2.Message pb = publishingMessage.toProtobuf(FAKE_NAMESPACE, mq);
assertEquals(apache.rocketmq.v2.Encoding.IDENTITY, pb.getSystemProperties().getBodyEncoding());
assertArrayEquals(body, pb.getBody().toByteArray());
}

@Test
public void testBodyCompressedWhenThresholdReached() throws IOException {
String topic = "testTopic";
byte[] body = new byte[4096];
Arrays.fill(body, (byte) 'x');
final Message message = new MessageBuilderImpl().setTopic(topic).setBody(body).build();
final PublishingSettings settings = fakeProducerSettings();
PublishingSettingsTestHelper.setCompressBodyThresholdBytes(settings, 1024);
final PublishingMessageImpl publishingMessage = new PublishingMessageImpl(message, settings, false);
final MessageQueueImpl mq = fakeMessageQueueImpl(topic);
final apache.rocketmq.v2.Message pb = publishingMessage.toProtobuf(FAKE_NAMESPACE, mq);
assertEquals(apache.rocketmq.v2.Encoding.GZIP, pb.getSystemProperties().getBodyEncoding());
final byte[] transportBody = pb.getBody().toByteArray();
assertTrue(transportBody.length < body.length);
assertArrayEquals(body, Utilities.decompressBytes(transportBody));
}

@Test
public void testBodyNotCompressedBelowThreshold() throws IOException {
String topic = "testTopic";
byte[] body = "foobar".getBytes(StandardCharsets.UTF_8);
final Message message = new MessageBuilderImpl().setTopic(topic).setBody(body).build();
final PublishingSettings settings = fakeProducerSettings();
PublishingSettingsTestHelper.setCompressBodyThresholdBytes(settings, 1024);
final PublishingMessageImpl publishingMessage = new PublishingMessageImpl(message, settings, false);
final MessageQueueImpl mq = fakeMessageQueueImpl(topic);
final apache.rocketmq.v2.Message pb = publishingMessage.toProtobuf(FAKE_NAMESPACE, mq);
assertEquals(apache.rocketmq.v2.Encoding.IDENTITY, pb.getSystemProperties().getBodyEncoding());
assertArrayEquals(body, pb.getBody().toByteArray());
}

@Test
public void testBodyCompressedAtExactThreshold() throws IOException {
String topic = "testTopic";
byte[] body = new byte[1024];
Arrays.fill(body, (byte) 'x');
final Message message = new MessageBuilderImpl().setTopic(topic).setBody(body).build();
final PublishingSettings settings = fakeProducerSettings();
PublishingSettingsTestHelper.setCompressBodyThresholdBytes(settings, 1024);
final PublishingMessageImpl publishingMessage = new PublishingMessageImpl(message, settings, false);
final MessageQueueImpl mq = fakeMessageQueueImpl(topic);
final apache.rocketmq.v2.Message pb = publishingMessage.toProtobuf(FAKE_NAMESPACE, mq);
assertEquals(apache.rocketmq.v2.Encoding.GZIP, pb.getSystemProperties().getBodyEncoding());
assertArrayEquals(body, Utilities.decompressBytes(pb.getBody().toByteArray()));
}

@Test
public void testBodyNotCompressedWhenGzipInflates() throws IOException {
String topic = "testTopic";
// Incompressible pseudo-random body: GZIP output is larger than the input,
// so the implementation must fall back to the identity encoding.
byte[] body = new byte[4096];
new Random(42).nextBytes(body);
final Message message = new MessageBuilderImpl().setTopic(topic).setBody(body).build();
final PublishingSettings settings = fakeProducerSettings();
PublishingSettingsTestHelper.setCompressBodyThresholdBytes(settings, 1024);
final PublishingMessageImpl publishingMessage = new PublishingMessageImpl(message, settings, false);
final MessageQueueImpl mq = fakeMessageQueueImpl(topic);
final apache.rocketmq.v2.Message pb = publishingMessage.toProtobuf(FAKE_NAMESPACE, mq);
assertEquals(apache.rocketmq.v2.Encoding.IDENTITY, pb.getSystemProperties().getBodyEncoding());
assertArrayEquals(body, pb.getBody().toByteArray());
}

@Test(expected = IOException.class)
public void testIncompressibleBodyExceedingMaxSizeRejected() throws IOException {
String topic = "testTopic";
// Larger than the default 4 MiB limit and incompressible: the transport body
// stays over the limit after the compression fallback, so it must be rejected.
byte[] body = new byte[5 * 1024 * 1024];
new Random(42).nextBytes(body);
final Message message = new MessageBuilderImpl().setTopic(topic).setBody(body).build();
final PublishingSettings settings = fakeProducerSettings();
PublishingSettingsTestHelper.setCompressBodyThresholdBytes(settings, 1024);
new PublishingMessageImpl(message, settings, false);
}

@Test
public void testCompressibleBodyOverMaxSizeAcceptedWhenCompressedFits() throws IOException {
String topic = "testTopic";
// Original body exceeds the 4 MiB limit, but the compressed body fits: the
// limit is validated against the transport (compressed) body by design.
byte[] body = new byte[5 * 1024 * 1024];
Arrays.fill(body, (byte) 'x');
final Message message = new MessageBuilderImpl().setTopic(topic).setBody(body).build();
final PublishingSettings settings = fakeProducerSettings();
PublishingSettingsTestHelper.setCompressBodyThresholdBytes(settings, 1024);
final PublishingMessageImpl publishingMessage = new PublishingMessageImpl(message, settings, false);
final MessageQueueImpl mq = fakeMessageQueueImpl(topic);
final apache.rocketmq.v2.Message pb = publishingMessage.toProtobuf(FAKE_NAMESPACE, mq);
assertEquals(apache.rocketmq.v2.Encoding.GZIP, pb.getSystemProperties().getBodyEncoding());
assertArrayEquals(body, Utilities.decompressBytes(pb.getBody().toByteArray()));
}
}