From d2a3238b0e975135c833169e84a5b7a7a04d6924 Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Thu, 9 May 2019 20:50:59 +0800 Subject: [PATCH 1/8] Support set message size --- *Motivation* Currently Pulsar only support 5MB size of messages.But there are many cases will use more than 5MB message to transfer. https://github.com/apache/pulsar/wiki/PIP-36%3A-Max-Message-Size *Modifications* - Add message size in protocol - Automaticlly adjust client message size by server --- .../pulsar/broker/ServiceConfiguration.java | 5 ++ .../apache/pulsar/PulsarBrokerStarter.java | 10 ++++ .../service/PulsarChannelInitializer.java | 5 +- .../pulsar/broker/service/ServerCnx.java | 4 +- .../apache/pulsar/client/impl/ClientCnx.java | 10 +++- .../pulsar/client/impl/ConsumerImpl.java | 2 +- .../pulsar/client/impl/ProducerImpl.java | 17 +++--- .../client/impl/PulsarChannelInitializer.java | 5 +- .../impl/conf/ClientConfigurationData.java | 1 + .../apache/pulsar/common/api/Commands.java | 3 +- .../pulsar/common/api/proto/PulsarApi.java | 57 +++++++++++++++++++ pulsar-common/src/main/proto/PulsarApi.proto | 1 + .../discovery/service/ServerConnection.java | 2 +- .../service/ServiceChannelInitializer.java | 6 +- .../service/server/ServiceConfig.java | 10 ++++ .../proxy/server/DirectProxyHandler.java | 5 +- .../proxy/server/ProxyConfiguration.java | 3 + .../pulsar/proxy/server/ProxyConnection.java | 3 +- .../server/ServiceChannelInitializer.java | 3 +- 19 files changed, 131 insertions(+), 21 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 5428d557a90e5..bc17f5e0710f5 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -489,6 +489,11 @@ public class ServiceConfiguration implements PulsarConfiguration { + " Using a value of 0, is disabling maxConsumersPerSubscription-limit check.") private int maxConsumersPerSubscription = 0; + @FieldContext(category = CATEGORY_SERVER, doc = "The size of messages.") + private int maxMessageSize = 5 * 1024 * 1024; + + private int maxFrameSize = 5242880; + /***** --- TLS --- ****/ @FieldContext( category = CATEGORY_TLS, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java index 3ef104f82fdd9..0b73a1042e968 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java @@ -45,6 +45,7 @@ import org.apache.bookkeeper.replication.AutoRecoveryMain; import org.apache.bookkeeper.stats.StatsProvider; import org.apache.bookkeeper.common.util.ReflectionUtils; +import org.apache.bookkeeper.util.DirectMemoryUtils; import org.apache.commons.configuration.ConfigurationException; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; @@ -139,6 +140,15 @@ private static class BrokerStarter { brokerConfig = loadConfig(starterArguments.brokerConfigFile); } + int maxFrameSize = brokerConfig.getMaxMessageSize() + (10 * 1024); + if (maxFrameSize < 0) { + throw new IllegalArgumentException("Max message size need smaller than 5233640 bytes"); + } + if (maxFrameSize > DirectMemoryUtils.maxDirectMemory()) { + throw new IllegalArgumentException("Max message size need smaller than jvm directMemory"); + } + brokerConfig.setMaxFrameSize(maxFrameSize); + // init functions worker if (starterArguments.runFunctionsWorker || brokerConfig.isFunctionsWorkerEnabled()) { WorkerConfig workerConfig; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java index ff24caa2ce797..f807a37ab29bf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java @@ -21,7 +21,6 @@ import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.common.api.ByteBufPair; -import org.apache.pulsar.common.api.PulsarDecoder; import org.apache.pulsar.common.util.NettySslContextBuilder; import io.netty.channel.ChannelInitializer; @@ -35,6 +34,7 @@ public class PulsarChannelInitializer extends ChannelInitializer private final PulsarService pulsar; private final boolean enableTls; private final NettySslContextBuilder sslCtxRefresher; + private final ServiceConfiguration brokerConf; /** * @@ -54,6 +54,7 @@ public PulsarChannelInitializer(PulsarService pulsar, boolean enableTLS) throws } else { this.sslCtxRefresher = null; } + this.brokerConf = pulsar.getConfiguration(); } @Override @@ -65,7 +66,7 @@ protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast("ByteBufPairEncoder", ByteBufPair.ENCODER); } - ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder(PulsarDecoder.MaxFrameSize, 0, 4, 0, 4)); + ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder(brokerConf.getMaxFrameSize(), 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ServerCnx(pulsar)); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 6fe386488848f..b5b6fe97b2d97 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -136,6 +136,7 @@ public class ServerCnx extends PulsarHandler { private boolean authenticateOriginalAuthData; private final boolean schemaValidationEnforced; private String authMethod = "none"; + private final int maxMessageSize; enum State { Start, Connected, Failed, Connecting @@ -156,6 +157,7 @@ public ServerCnx(PulsarService pulsar) { this.proxyRoles = service.pulsar().getConfiguration().getProxyRoles(); this.authenticateOriginalAuthData = service.pulsar().getConfiguration().isAuthenticateOriginalAuthData(); this.schemaValidationEnforced = pulsar.getConfiguration().isSchemaValidationEnforced(); + this.maxMessageSize = pulsar.getConfiguration().getMaxMessageSize(); } @Override @@ -455,7 +457,7 @@ private String getOriginalPrincipal(String originalAuthData, String originalAuth // complete the connect and sent newConnected command private void completeConnect(int clientProtoVersion, String clientVersion) { - ctx.writeAndFlush(Commands.newConnected(clientProtoVersion)); + ctx.writeAndFlush(Commands.newConnected(clientProtoVersion, maxMessageSize)); state = State.Connected; remoteEndpointProtocolVersion = clientProtoVersion; if (isNotBlank(clientVersion) && !clientVersion.contains(" ") /* ignore default version: pulsar client */) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index 1a88afd0c4bf2..7eaf10a67e2e6 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -30,6 +30,7 @@ import io.netty.channel.ChannelHandlerContext; import io.netty.channel.EventLoopGroup; import io.netty.channel.unix.Errors.NativeIoException; +import io.netty.handler.codec.LengthFieldBasedFrameDecoder; import io.netty.handler.ssl.SslHandler; import io.netty.util.concurrent.Promise; @@ -117,6 +118,10 @@ public class ClientCnx extends PulsarHandler { .newUpdater(ClientCnx.class, "numberOfRejectRequests"); @SuppressWarnings("unused") private volatile int numberOfRejectRequests = 0; + + protected static final AtomicIntegerFieldUpdater NUMBER_OF_MAX_MESSAGE_SIZE = AtomicIntegerFieldUpdater + .newUpdater(ClientCnx.class, "maxMessageSize"); + private volatile int maxMessageSize = 0; private final int maxNumberOfRejectedRequestPerConnection; private final int rejectedRequestResetTimeSec = 60; private final int protocolVersion; @@ -275,7 +280,10 @@ protected void handleConnected(CommandConnected connected) { } checkArgument(state == State.SentConnectFrame || state == State.Connecting); - + int maxFrameSize = connected.getMaxMessageSize() - 10 * 1024; + NUMBER_OF_MAX_MESSAGE_SIZE.compareAndSet(this, 0, maxFrameSize); + ctx.pipeline() + .replace("defaultFrameDecoder", "frameDecoder", new LengthFieldBasedFrameDecoder(maxFrameSize, 0, 4, 0, 4)); if (log.isDebugEnabled()) { log.debug("{} Connection is ready", ctx.channel()); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 8909989e5ada3..cee610f6c0f33 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -1141,7 +1141,7 @@ private ByteBuf uncompressPayloadIfNeeded(MessageIdData messageId, MessageMetada CompressionCodec codec = CompressionCodecProvider.getCompressionCodec(compressionType); int uncompressedSize = msgMetadata.getUncompressedSize(); int payloadSize = payload.readableBytes(); - if (payloadSize > PulsarDecoder.MaxMessageSize) { + if (payloadSize > ClientCnx.NUMBER_OF_MAX_MESSAGE_SIZE.get(currentCnx)) { // payload size is itself corrupted since it cannot be bigger than the MaxMessageSize log.error("[{}][{}] Got corrupted payload message size {} at {}", topic, subscription, payloadSize, messageId); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index b4a0993b49c8c..2cc978bbfcd04 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -312,15 +312,15 @@ public void sendAsync(Message message, SendCallback callback) { // validate msg-size (For batching this will be check at the batch completion size) int compressedSize = compressedPayload.readableBytes(); - - if (compressedSize > PulsarDecoder.MaxMessageSize) { + int maxMessageSize = ClientCnx.NUMBER_OF_MAX_MESSAGE_SIZE.get(this.cnx()); + if (compressedSize > maxMessageSize) { compressedPayload.release(); String compressedStr = (!isBatchMessagingEnabled() && conf.getCompressionType() != CompressionType.NONE) - ? "Compressed" - : ""; + ? "Compressed" + : ""; PulsarClientException.InvalidMessageException invalidMessageException = new PulsarClientException.InvalidMessageException( - format("%s Message payload size %d cannot exceed %d bytes", compressedStr, compressedSize, - PulsarDecoder.MaxMessageSize)); + format("%s Message payload size %d cannot exceed %d bytes", compressedStr, compressedSize, + maxMessageSize)); callback.sendComplete(invalidMessageException); return; } @@ -1306,12 +1306,13 @@ private void batchMessageAndSend() { op = OpSendMsg.create(batchMessageContainer.messages, cmd, sequenceId, batchMessageContainer.firstCallback); - if (encryptedPayload.readableBytes() > PulsarDecoder.MaxMessageSize) { + int maxMessageSize = ClientCnx.NUMBER_OF_MAX_MESSAGE_SIZE.get(this.cnx()); + if (encryptedPayload.readableBytes() > maxMessageSize) { cmd.release(); semaphore.release(numMessagesInBatch); if (op != null) { op.callback.sendComplete(new PulsarClientException.InvalidMessageException( - "Message size is bigger than " + PulsarDecoder.MaxMessageSize + " bytes")); + "Message size is bigger than " + maxMessageSize + " bytes")); } return; } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java index 494cb1283264b..6688bca7d0c18 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java @@ -38,6 +38,7 @@ public class PulsarChannelInitializer extends ChannelInitializer private final Supplier clientCnxSupplier; private final SslContext sslCtx; + private final ClientConfigurationData conf; public PulsarChannelInitializer(ClientConfigurationData conf, Supplier clientCnxSupplier) throws Exception { @@ -57,6 +58,7 @@ public PulsarChannelInitializer(ClientConfigurationData conf, Supplier private final DiscoveryService discoveryService; private final boolean enableTls; private final NettySslContextBuilder sslCtxRefresher; + private final ServiceConfig brokerConf; public ServiceChannelInitializer(DiscoveryService discoveryService, ServiceConfig serviceConfig, boolean e) throws Exception { @@ -53,6 +55,7 @@ public ServiceChannelInitializer(DiscoveryService discoveryService, ServiceConfi } else { this.sslCtxRefresher = null; } + this.brokerConf = serviceConfig; } @Override @@ -63,7 +66,8 @@ protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast(TLS_HANDLER, sslContext.newHandler(ch.alloc())); } } - ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder(PulsarDecoder.MaxFrameSize, 0, 4, 0, 4)); + ch.pipeline().addLast("frameDecoder", + new LengthFieldBasedFrameDecoder(brokerConf.getMaxMessageSize() + 10 * 1024, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ServerConnection(discoveryService)); } } diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java index 5f89527a8d241..bba71dfc90185 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java @@ -75,6 +75,16 @@ public class ServiceConfig implements PulsarConfiguration { // Authorization provider fully qualified class-name private String authorizationProvider = PulsarAuthorizationProvider.class.getName(); + public int getMaxMessageSize() { + return maxMessageSize; + } + + public void setMaxMessageSize(int maxMessageSize) { + this.maxMessageSize = maxMessageSize; + } + + private int maxMessageSize = 5242880; + /***** --- TLS --- ****/ @Deprecated private boolean tlsEnabled = false; diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index ffa4c2c838cb9..530b2adc1c3ae 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -99,7 +99,8 @@ protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast(TLS_HANDLER, sslCtx.newHandler(ch.alloc())); } ch.pipeline().addLast("frameDecoder", - new LengthFieldBasedFrameDecoder(PulsarDecoder.MaxFrameSize, 0, 4, 0, 4)); + new LengthFieldBasedFrameDecoder(config.getMaxMessagesSize() + 10 * 1024, 0, 4, 0, + 4)); ch.pipeline().addLast("proxyOutboundHandler", new ProxyBackendHandler(config, protocolVersion)); } }); @@ -259,7 +260,7 @@ protected void handleConnected(CommandConnected connected) { state = BackendState.HandshakeCompleted; - inboundChannel.writeAndFlush(Commands.newConnected(connected.getProtocolVersion())).addListener(future -> { + inboundChannel.writeAndFlush(Commands.newConnected(connected.getProtocolVersion(), config.getMaxMessagesSize())).addListener(future -> { if (log.isDebugEnabled()) { log.debug("[{}] [{}] Removing decoder from pipeline", inboundChannel, outboundChannel); } diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java index 1ca1159c772c3..e728f9045f51e 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java @@ -111,6 +111,9 @@ public class ProxyConfiguration implements PulsarConfiguration { ) private String brokerWebServiceURLTLS; + @FieldContext(category = CATEGORY_BROKER_DISCOVERY, doc = "") + private int maxMessagesSize = 5233640; + @FieldContext( category = CATEGORY_BROKER_DISCOVERY, doc = "The web service url points to the function worker cluster." diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index 18e4e5d750822..a9c86233fd766 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -218,7 +218,8 @@ private void completeConnect() { // partitions metadata lookups state = State.ProxyLookupRequests; lookupProxyHandler = new LookupProxyHandler(service, this); - ctx.writeAndFlush(Commands.newConnected(protocolVersionToAdvertise)); + ctx.writeAndFlush( + Commands.newConnected(protocolVersionToAdvertise, service.getConfiguration().getMaxMessagesSize())); } } diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java index b03afbfdf4460..8ac139a5e4acf 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java @@ -85,7 +85,8 @@ protected void initChannel(SocketChannel ch) throws Exception { } } - ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder(PulsarDecoder.MaxFrameSize, 0, 4, 0, 4)); + ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( + proxyService.getConfiguration().getMaxMessagesSize() + 10 * 1024, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ProxyConnection(proxyService, clientSslCtxRefresher == null ? null : clientSslCtxRefresher.get())); } From 6bbd00f7809fd84877cb0484d6b0a51a2cba5d59 Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Fri, 10 May 2019 18:29:26 +0800 Subject: [PATCH 2/8] Use `maxMessageSize` to set `nettyFrameSize` in bookie client --- *Motivation* When broker specify a `maxMessageSize` bookie should accept this value as `nettyFrameSize` *Modifications* - Use `cnx().getMaxMessageSize` - Discovery service only redirect so use the constant value `5 * 1024 * 1024` as message size - Put `MAX_METADATA_SIZE` as constant value in `InternalConfigurationData` --- .../pulsar/broker/ServiceConfiguration.java | 8 +-- .../apache/pulsar/PulsarBrokerStarter.java | 9 ++-- .../broker/BookKeeperClientFactoryImpl.java | 2 + .../service/PulsarChannelInitializer.java | 4 +- .../apache/pulsar/client/impl/ClientCnx.java | 15 +++--- .../pulsar/client/impl/ConsumerImpl.java | 5 +- .../pulsar/client/impl/ProducerImpl.java | 11 ++--- .../conf/InternalConfigurationData.java | 1 + pulsar-common/src/main/proto/PulsarApi.proto | 2 +- .../discovery/service/ServerConnection.java | 3 +- .../service/ServiceChannelInitializer.java | 8 ++- .../service/server/ServiceConfig.java | 17 ++----- .../proxy/server/DirectProxyHandler.java | 49 ++++++++++++------- .../server/ServiceChannelInitializer.java | 8 +-- 14 files changed, 73 insertions(+), 69 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index bc17f5e0710f5..d987ba9e6e900 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -33,6 +33,7 @@ import lombok.Setter; import org.apache.bookkeeper.client.api.DigestType; import org.apache.pulsar.broker.authorization.PulsarAuthorizationProvider; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.configuration.Category; import org.apache.pulsar.common.configuration.FieldContext; import org.apache.pulsar.common.configuration.PulsarConfiguration; @@ -489,11 +490,12 @@ public class ServiceConfiguration implements PulsarConfiguration { + " Using a value of 0, is disabling maxConsumersPerSubscription-limit check.") private int maxConsumersPerSubscription = 0; - @FieldContext(category = CATEGORY_SERVER, doc = "The size of messages.") + @FieldContext( + category = CATEGORY_SERVER, + doc = "Max size of messages.", + maxValue = Integer.MAX_VALUE - InternalConfigurationData.MESSAGE_META_SIZE) private int maxMessageSize = 5 * 1024 * 1024; - private int maxFrameSize = 5242880; - /***** --- TLS --- ****/ @FieldContext( category = CATEGORY_TLS, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java index 0b73a1042e968..4b8ae506a71a1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java @@ -50,6 +50,7 @@ import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.ServiceConfigurationUtils; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.functions.worker.WorkerConfig; import org.apache.pulsar.functions.worker.WorkerService; import org.slf4j.Logger; @@ -140,14 +141,10 @@ private static class BrokerStarter { brokerConfig = loadConfig(starterArguments.brokerConfigFile); } - int maxFrameSize = brokerConfig.getMaxMessageSize() + (10 * 1024); - if (maxFrameSize < 0) { - throw new IllegalArgumentException("Max message size need smaller than 5233640 bytes"); - } - if (maxFrameSize > DirectMemoryUtils.maxDirectMemory()) { + int maxFrameSize = brokerConfig.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE; + if (maxFrameSize >= DirectMemoryUtils.maxDirectMemory()) { throw new IllegalArgumentException("Max message size need smaller than jvm directMemory"); } - brokerConfig.setMaxFrameSize(maxFrameSize); // init functions worker if (starterArguments.runFunctionsWorker || brokerConfig.isFunctionsWorkerEnabled()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java index f15b9a3f8d70b..b80fc786cc4be 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java @@ -28,6 +28,7 @@ import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy; import org.apache.bookkeeper.client.RegionAwareEnsemblePlacementPolicy; import org.apache.bookkeeper.conf.ClientConfiguration; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.zookeeper.ZkBookieRackAffinityMapping; import org.apache.pulsar.zookeeper.ZkIsolatedBookieEnsemblePlacementPolicy; import org.apache.pulsar.zookeeper.ZooKeeperCache; @@ -57,6 +58,7 @@ public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient) throws I bkConf.setUseV2WireProtocol(conf.isBookkeeperUseV2WireProtocol()); bkConf.setEnableDigestTypeAutodetection(true); bkConf.setStickyReadsEnabled(conf.isBookkeeperEnableStickyReads()); + bkConf.setNettyMaxFrameSizeBytes(conf.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE); bkConf.setAllocatorPoolingPolicy(PoolingPolicy.UnpooledHeap); if (conf.isBookkeeperClientHealthCheckEnabled()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java index f807a37ab29bf..6d622b5e019c3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java @@ -21,6 +21,7 @@ import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.common.api.ByteBufPair; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.util.NettySslContextBuilder; import io.netty.channel.ChannelInitializer; @@ -66,7 +67,8 @@ protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast("ByteBufPairEncoder", ByteBufPair.ENCODER); } - ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder(brokerConf.getMaxFrameSize(), 0, 4, 0, 4)); + ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( + brokerConf.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ServerCnx(pulsar)); } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index 7eaf10a67e2e6..5866cf2be67f9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -49,6 +49,7 @@ import javax.net.ssl.SSLSession; +import lombok.Getter; import org.apache.commons.lang3.tuple.Pair; import org.apache.http.conn.ssl.DefaultHostnameVerifier; import org.apache.pulsar.PulsarVersion; @@ -80,6 +81,7 @@ import org.apache.pulsar.common.api.proto.PulsarApi.CommandSuccess; import org.apache.pulsar.common.api.proto.PulsarApi.MessageIdData; import org.apache.pulsar.common.api.proto.PulsarApi.ServerError; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaInfoUtil; import org.apache.pulsar.common.util.collections.ConcurrentLongHashMap; @@ -119,9 +121,9 @@ public class ClientCnx extends PulsarHandler { @SuppressWarnings("unused") private volatile int numberOfRejectRequests = 0; - protected static final AtomicIntegerFieldUpdater NUMBER_OF_MAX_MESSAGE_SIZE = AtomicIntegerFieldUpdater - .newUpdater(ClientCnx.class, "maxMessageSize"); - private volatile int maxMessageSize = 0; + @Getter + private int maxMessageSize = 0; + private final int maxNumberOfRejectedRequestPerConnection; private final int rejectedRequestResetTimeSec = 60; private final int protocolVersion; @@ -280,10 +282,9 @@ protected void handleConnected(CommandConnected connected) { } checkArgument(state == State.SentConnectFrame || state == State.Connecting); - int maxFrameSize = connected.getMaxMessageSize() - 10 * 1024; - NUMBER_OF_MAX_MESSAGE_SIZE.compareAndSet(this, 0, maxFrameSize); - ctx.pipeline() - .replace("defaultFrameDecoder", "frameDecoder", new LengthFieldBasedFrameDecoder(maxFrameSize, 0, 4, 0, 4)); + this.maxMessageSize = connected.getMaxMessageSize(); + ctx.pipeline().replace("defaultFrameDecoder", "frameDecoder", new LengthFieldBasedFrameDecoder( + connected.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, 4)); if (log.isDebugEnabled()) { log.debug("{} Connection is ready", ctx.channel()); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index cee610f6c0f33..9946c48ee3122 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -66,7 +66,6 @@ import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.api.EncryptionContext; import org.apache.pulsar.common.api.EncryptionContext.EncryptionKey; -import org.apache.pulsar.common.api.PulsarDecoder; import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.AckType; import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.ValidationError; @@ -1141,10 +1140,10 @@ private ByteBuf uncompressPayloadIfNeeded(MessageIdData messageId, MessageMetada CompressionCodec codec = CompressionCodecProvider.getCompressionCodec(compressionType); int uncompressedSize = msgMetadata.getUncompressedSize(); int payloadSize = payload.readableBytes(); - if (payloadSize > ClientCnx.NUMBER_OF_MAX_MESSAGE_SIZE.get(currentCnx)) { + if (payloadSize > cnx().getMaxMessageSize()) { // payload size is itself corrupted since it cannot be bigger than the MaxMessageSize log.error("[{}][{}] Got corrupted payload message size {} at {}", topic, subscription, payloadSize, - messageId); + messageId); discardCorruptedMessage(messageId, currentCnx, ValidationError.UncompressedSizeCorruption); return null; } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 2cc978bbfcd04..831eba2f7b2a3 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -61,7 +61,6 @@ import org.apache.pulsar.common.api.ByteBufPair; import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.api.Commands.ChecksumType; -import org.apache.pulsar.common.api.PulsarDecoder; import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; import org.apache.pulsar.common.api.proto.PulsarApi.ProtocolVersion; import org.apache.pulsar.common.compression.CompressionCodec; @@ -312,15 +311,14 @@ public void sendAsync(Message message, SendCallback callback) { // validate msg-size (For batching this will be check at the batch completion size) int compressedSize = compressedPayload.readableBytes(); - int maxMessageSize = ClientCnx.NUMBER_OF_MAX_MESSAGE_SIZE.get(this.cnx()); - if (compressedSize > maxMessageSize) { + if (compressedSize > cnx().getMaxMessageSize()) { compressedPayload.release(); String compressedStr = (!isBatchMessagingEnabled() && conf.getCompressionType() != CompressionType.NONE) ? "Compressed" : ""; PulsarClientException.InvalidMessageException invalidMessageException = new PulsarClientException.InvalidMessageException( format("%s Message payload size %d cannot exceed %d bytes", compressedStr, compressedSize, - maxMessageSize)); + cnx().getMaxMessageSize())); callback.sendComplete(invalidMessageException); return; } @@ -1306,13 +1304,12 @@ private void batchMessageAndSend() { op = OpSendMsg.create(batchMessageContainer.messages, cmd, sequenceId, batchMessageContainer.firstCallback); - int maxMessageSize = ClientCnx.NUMBER_OF_MAX_MESSAGE_SIZE.get(this.cnx()); - if (encryptedPayload.readableBytes() > maxMessageSize) { + if (encryptedPayload.readableBytes() > cnx().getMaxMessageSize()) { cmd.release(); semaphore.release(numMessagesInBatch); if (op != null) { op.callback.sendComplete(new PulsarClientException.InvalidMessageException( - "Message size is bigger than " + maxMessageSize + " bytes")); + "Message size is bigger than " + cnx().getMaxMessageSize() + " bytes")); } return; } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java b/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java index 558c1aa8da8a5..45d9c110796c4 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java @@ -23,6 +23,7 @@ public class InternalConfigurationData { + public final static int MESSAGE_META_SIZE = 10 * 1024; private String zookeeperServers; private String configurationStoreServers; private String ledgersRootPath; diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index 02a431c4279f2..11d69bee21c48 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -213,7 +213,7 @@ message CommandConnect { message CommandConnected { required string server_version = 1; optional int32 protocol_version = 2 [default = 0]; - optional int32 max_message_size = 3 [default = 5242880]; + optional int32 max_message_size = 3; } message CommandAuthResponse { diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java index 39b9ea32354b3..a35c6aadedd62 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java @@ -36,6 +36,7 @@ import org.apache.pulsar.common.api.proto.PulsarApi.CommandPartitionedTopicMetadata; import org.apache.pulsar.common.api.proto.PulsarApi.ServerError; import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.discovery.service.server.ServiceConfig; import org.apache.pulsar.policies.data.loadbalancer.LoadManagerReport; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -101,7 +102,7 @@ protected void handleConnect(CommandConnect connect) { return; } } - ctx.writeAndFlush(Commands.newConnected(connect.getProtocolVersion(), service.getConfiguration().getMaxMessageSize())); + ctx.writeAndFlush(Commands.newConnected(connect.getProtocolVersion(), ServiceConfig.MAX_MESSAGE_SIZE)); state = State.Connected; remoteEndpointProtocolVersion = connect.getProtocolVersion(); } diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java index 8af24645aa06c..0224d91a15613 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java @@ -18,10 +18,8 @@ */ package org.apache.pulsar.discovery.service; -import org.apache.pulsar.broker.ServiceConfiguration; -import org.apache.pulsar.common.api.PulsarDecoder; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.util.NettySslContextBuilder; -import org.apache.pulsar.common.util.SslContextAutoRefreshBuilder; import org.apache.pulsar.discovery.service.server.ServiceConfig; import io.netty.channel.ChannelInitializer; @@ -66,8 +64,8 @@ protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast(TLS_HANDLER, sslContext.newHandler(ch.alloc())); } } - ch.pipeline().addLast("frameDecoder", - new LengthFieldBasedFrameDecoder(brokerConf.getMaxMessageSize() + 10 * 1024, 0, 4, 0, 4)); + ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( + ServiceConfig.MAX_MESSAGE_SIZE + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ServerConnection(discoveryService)); } } diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java index bba71dfc90185..5aedb59735221 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java @@ -1,5 +1,4 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one +/** * 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 @@ -34,6 +33,10 @@ */ public class ServiceConfig implements PulsarConfiguration { + // Discovery service doesn't send any messages except Command connected. + // So it's ok use a default value. + public final static int MAX_MESSAGE_SIZE = 5 * 1024 * 1024; + // Local-Zookeeper quorum connection string private String zookeeperServers; // Global-Zookeeper quorum connection string @@ -75,16 +78,6 @@ public class ServiceConfig implements PulsarConfiguration { // Authorization provider fully qualified class-name private String authorizationProvider = PulsarAuthorizationProvider.class.getName(); - public int getMaxMessageSize() { - return maxMessageSize; - } - - public void setMaxMessageSize(int maxMessageSize) { - this.maxMessageSize = maxMessageSize; - } - - private int maxMessageSize = 5242880; - /***** --- TLS --- ****/ @Deprecated private boolean tlsEnabled = false; diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index 530b2adc1c3ae..76d3f7906f7d7 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -21,7 +21,6 @@ import static com.google.common.base.Preconditions.checkArgument; import static com.google.common.base.Preconditions.checkState; -import static java.nio.charset.StandardCharsets.UTF_8; import java.net.URI; import java.net.URISyntaxException; @@ -39,6 +38,7 @@ import org.apache.pulsar.common.api.PulsarDecoder; import org.apache.pulsar.common.api.proto.PulsarApi.CommandAuthChallenge; import org.apache.pulsar.common.api.proto.PulsarApi.CommandConnected; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -260,24 +260,35 @@ protected void handleConnected(CommandConnected connected) { state = BackendState.HandshakeCompleted; - inboundChannel.writeAndFlush(Commands.newConnected(connected.getProtocolVersion(), config.getMaxMessagesSize())).addListener(future -> { - if (log.isDebugEnabled()) { - log.debug("[{}] [{}] Removing decoder from pipeline", inboundChannel, outboundChannel); - } - if (ProxyService.proxyLogLevel == 0) { - // direct tcp proxy - inboundChannel.pipeline().remove("frameDecoder"); - outboundChannel.pipeline().remove("frameDecoder"); - } else { - // Enable parsing feature, proxyLogLevel(1 or 2) - // Add parser handler - inboundChannel.pipeline().addBefore("handler" , "inboundParser" , new ParserProxyHandler(inboundChannel , ParserProxyHandler.FRONTEND_CONN)); - outboundChannel.pipeline().addBefore("proxyOutboundHandler" , "outboundParser" , new ParserProxyHandler(outboundChannel , ParserProxyHandler.BACKEND_CONN)); - } - // Start reading from both connections - inboundChannel.read(); - outboundChannel.read(); - }); + inboundChannel + .writeAndFlush(Commands.newConnected(connected.getProtocolVersion(), connected.getMaxMessageSize())) + .addListener(future -> { + if (log.isDebugEnabled()) { + log.debug("[{}] [{}] Removing decoder from pipeline", inboundChannel, outboundChannel); + } + if (ProxyService.proxyLogLevel == 0) { + // direct tcp proxy + inboundChannel.pipeline().remove("defaultFrameDecoder"); + outboundChannel.pipeline().remove("defaultFrameDecoder"); + } else { + // Enable parsing feature, proxyLogLevel(1 or 2) + // Add parser handler + inboundChannel.pipeline().replace("defaultFrameDecoder", "frameDecoder", + new LengthFieldBasedFrameDecoder( + connected.getMaxMessageSize() + + InternalConfigurationData.MESSAGE_META_SIZE, + 0, 4, 0, 4)); + inboundChannel.pipeline().addBefore("handler", "inboundParser", + new ParserProxyHandler(inboundChannel, + ParserProxyHandler.FRONTEND_CONN)); + outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", + new ParserProxyHandler(outboundChannel, + ParserProxyHandler.BACKEND_CONN)); + } + // Start reading from both connections + inboundChannel.read(); + outboundChannel.read(); + }); } @Override diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java index 8ac139a5e4acf..704518885b0b1 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java @@ -22,10 +22,9 @@ import org.apache.pulsar.client.api.AuthenticationDataProvider; import org.apache.pulsar.client.api.AuthenticationFactory; -import org.apache.pulsar.common.api.PulsarDecoder; +import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.util.ClientSslContextRefresher; import org.apache.pulsar.common.util.NettySslContextBuilder; -import org.apache.pulsar.common.util.SslContextAutoRefreshBuilder; import io.netty.channel.ChannelInitializer; import io.netty.channel.socket.SocketChannel; @@ -85,8 +84,9 @@ protected void initChannel(SocketChannel ch) throws Exception { } } - ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( - proxyService.getConfiguration().getMaxMessagesSize() + 10 * 1024, 0, 4, 0, 4)); + ch.pipeline().addLast("defaultFrameDecoder", new LengthFieldBasedFrameDecoder( + proxyService.getConfiguration().getMaxMessagesSize() + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, + 4)); ch.pipeline().addLast("handler", new ProxyConnection(proxyService, clientSslCtxRefresher == null ? null : clientSslCtxRefresher.get())); } From 53f8ae2bc3c84a3676a22ecf6475f570e32b1fee Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Mon, 13 May 2019 16:49:01 +0800 Subject: [PATCH 3/8] Use `Commands` to store message setting --- *Modifications* - use `Commands` to store default `MAX_MESSAGE_SIZE` and `MESSAGE_SIZE_FRAME_PADDING` - replace `LengthFieldBasedFrameDecoder` when has set message size - replace `PulsarDecoder.MaxMessageSize` --- .../pulsar/broker/ServiceConfiguration.java | 5 +- .../apache/pulsar/PulsarBrokerStarter.java | 3 +- .../broker/BookKeeperClientFactoryImpl.java | 3 +- .../service/PulsarChannelInitializer.java | 3 +- .../api/SimpleProducerConsumerTest.java | 14 ++-- .../api/v1/V1_ProducerConsumerTest.java | 5 +- .../pulsar/client/impl/MessageParserTest.java | 5 +- .../apache/pulsar/client/impl/ClientCnx.java | 15 ++-- .../client/impl/PulsarChannelInitializer.java | 6 +- .../impl/conf/ClientConfigurationData.java | 1 - .../apache/pulsar/common/api/Commands.java | 13 +++- .../pulsar/common/api/PulsarDecoder.java | 5 -- .../pulsar/common/api/raw/MessageParser.java | 8 +- .../discovery/service/ServerConnection.java | 2 +- .../service/ServiceChannelInitializer.java | 6 +- .../service/server/ServiceConfig.java | 4 - .../proxy/server/DirectProxyHandler.java | 74 ++++++++++++------- .../proxy/server/ParserProxyHandler.java | 9 ++- .../proxy/server/ProxyConfiguration.java | 3 - .../pulsar/proxy/server/ProxyConnection.java | 3 +- .../server/ServiceChannelInitializer.java | 7 +- .../sql/presto/PulsarConnectorConfig.java | 12 +++ .../pulsar/sql/presto/PulsarRecordCursor.java | 3 +- 23 files changed, 127 insertions(+), 82 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index d987ba9e6e900..726cba5f031ed 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -33,6 +33,7 @@ import lombok.Setter; import org.apache.bookkeeper.client.api.DigestType; import org.apache.pulsar.broker.authorization.PulsarAuthorizationProvider; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.configuration.Category; import org.apache.pulsar.common.configuration.FieldContext; @@ -493,8 +494,8 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext( category = CATEGORY_SERVER, doc = "Max size of messages.", - maxValue = Integer.MAX_VALUE - InternalConfigurationData.MESSAGE_META_SIZE) - private int maxMessageSize = 5 * 1024 * 1024; + maxValue = Integer.MAX_VALUE - Commands.MESSAGE_SIZE_FRAME_PADDING) + public int maxMessageSize = Commands.DEFAULT_MAX_MESSAGE_SIZE; /***** --- TLS --- ****/ @FieldContext( diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java index 4b8ae506a71a1..8f291ea31f4ef 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarBrokerStarter.java @@ -50,6 +50,7 @@ import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.ServiceConfigurationUtils; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.functions.worker.WorkerConfig; import org.apache.pulsar.functions.worker.WorkerService; @@ -141,7 +142,7 @@ private static class BrokerStarter { brokerConfig = loadConfig(starterArguments.brokerConfigFile); } - int maxFrameSize = brokerConfig.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE; + int maxFrameSize = brokerConfig.getMaxMessageSize() + Commands.MESSAGE_SIZE_FRAME_PADDING; if (maxFrameSize >= DirectMemoryUtils.maxDirectMemory()) { throw new IllegalArgumentException("Max message size need smaller than jvm directMemory"); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java index b80fc786cc4be..65207e6c67829 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java @@ -28,6 +28,7 @@ import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy; import org.apache.bookkeeper.client.RegionAwareEnsemblePlacementPolicy; import org.apache.bookkeeper.conf.ClientConfiguration; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.zookeeper.ZkBookieRackAffinityMapping; import org.apache.pulsar.zookeeper.ZkIsolatedBookieEnsemblePlacementPolicy; @@ -58,7 +59,7 @@ public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient) throws I bkConf.setUseV2WireProtocol(conf.isBookkeeperUseV2WireProtocol()); bkConf.setEnableDigestTypeAutodetection(true); bkConf.setStickyReadsEnabled(conf.isBookkeeperEnableStickyReads()); - bkConf.setNettyMaxFrameSizeBytes(conf.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE); + bkConf.setNettyMaxFrameSizeBytes(conf.getMaxMessageSize() + Commands.MESSAGE_SIZE_FRAME_PADDING); bkConf.setAllocatorPoolingPolicy(PoolingPolicy.UnpooledHeap); if (conf.isBookkeeperClientHealthCheckEnabled()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java index 6d622b5e019c3..f3e8214ffe1c9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/PulsarChannelInitializer.java @@ -21,6 +21,7 @@ import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.common.api.ByteBufPair; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.util.NettySslContextBuilder; @@ -68,7 +69,7 @@ protected void initChannel(SocketChannel ch) throws Exception { } ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( - brokerConf.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, 4)); + brokerConf.getMaxMessageSize() + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ServerCnx(pulsar)); } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index 7192e7fa6dd6a..32c6fa8f5ecb3 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -634,10 +634,10 @@ public void testSendBigMessageSize() throws Exception { // Messages are allowed up to MaxMessageSize - producer.newMessage().value(new byte[PulsarDecoder.MaxMessageSize]); + producer.newMessage().value(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE]); try { - producer.send(new byte[PulsarDecoder.MaxMessageSize + 1]); + producer.send(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE + 1]); fail("Should have thrown exception"); } catch (PulsarClientException.InvalidMessageException e) { // OK @@ -671,7 +671,7 @@ public void testSendBigMessageSizeButCompressed() throws Exception { .messageRoutingMode(MessageRoutingMode.SinglePartition) .compressionType(CompressionType.LZ4) .create(); - producer.send(new byte[PulsarDecoder.MaxMessageSize + 1]); + producer.send(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE + 1]); producer.close(); // (b) batch-msg with compression @@ -680,7 +680,7 @@ public void testSendBigMessageSizeButCompressed() throws Exception { .messageRoutingMode(MessageRoutingMode.SinglePartition) .compressionType(CompressionType.LZ4) .create(); - producer.send(new byte[PulsarDecoder.MaxMessageSize + 1]); + producer.send(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE + 1]); producer.close(); // (c) non-batch msg without compression @@ -690,7 +690,7 @@ public void testSendBigMessageSizeButCompressed() throws Exception { .compressionType(CompressionType.NONE) .create(); try { - producer.send(new byte[PulsarDecoder.MaxMessageSize + 1]); + producer.send(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE + 1]); fail("Should have thrown exception"); } catch (PulsarClientException.InvalidMessageException e) { // OK @@ -704,7 +704,7 @@ public void testSendBigMessageSizeButCompressed() throws Exception { .messageRoutingMode(MessageRoutingMode.SinglePartition) .compressionType(CompressionType.LZ4).create(); Consumer consumer = pulsarClient.newConsumer().topic(topic).subscriptionName("sub1").subscribe(); - byte[] content = new byte[PulsarDecoder.MaxMessageSize + 10]; + byte[] content = new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE + 10]; producer.send(content); assertEquals(consumer.receive().getData(), content); producer.close(); @@ -716,7 +716,7 @@ public void testSendBigMessageSizeButCompressed() throws Exception { .compressionType(CompressionType.NONE) .create(); try { - producer.send(new byte[PulsarDecoder.MaxMessageSize + 1]); + producer.send(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE + 1]); fail("Should have thrown exception"); } catch (PulsarClientException.InvalidMessageException e) { // OK diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java index fecef7946d5c1..2040979f36954 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java @@ -73,6 +73,7 @@ import org.apache.pulsar.client.impl.ConsumerImpl; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.TypedMessageBuilderImpl; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.api.PulsarDecoder; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.FutureUtil; @@ -607,10 +608,10 @@ public void testSendBigMessageSize() throws Exception { Producer producer = pulsarClient.newProducer().topic(topic).create(); // Messages are allowed up to MaxMessageSize - producer.newMessage().value(new byte[PulsarDecoder.MaxMessageSize]); + producer.newMessage().value(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE]); try { - producer.send(new byte[PulsarDecoder.MaxMessageSize + 1]); + producer.send(new byte[Commands.DEFAULT_MAX_MESSAGE_SIZE + 1]); fail("Should have thrown exception"); } catch (PulsarClientException.InvalidMessageException e) { // OK diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageParserTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageParserTest.java index 5ec6332a0ceef..1e0c45c2ac29c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageParserTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageParserTest.java @@ -34,6 +34,7 @@ import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.api.raw.MessageParser; import org.apache.pulsar.common.api.raw.RawMessage; import org.apache.pulsar.common.naming.TopicName; @@ -91,7 +92,7 @@ public void testWithoutBatches() throws Exception { MessageParser.parseMessage(topicName, entry.getLedgerId(), entry.getEntryId(), entry.getDataBuffer(), (message) -> { messages.add(message); - }); + }, Commands.DEFAULT_MAX_MESSAGE_SIZE); } finally { entry.release(); } @@ -133,7 +134,7 @@ public void testWithBatches() throws Exception { MessageParser.parseMessage(topicName, entry.getLedgerId(), entry.getEntryId(), entry.getDataBuffer(), (message) -> { messages.add(message); - }); + }, Commands.DEFAULT_MAX_MESSAGE_SIZE); } finally { entry.release(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index 5866cf2be67f9..b3ef8b2be35cd 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -81,7 +81,6 @@ import org.apache.pulsar.common.api.proto.PulsarApi.CommandSuccess; import org.apache.pulsar.common.api.proto.PulsarApi.MessageIdData; import org.apache.pulsar.common.api.proto.PulsarApi.ServerError; -import org.apache.pulsar.common.conf.InternalConfigurationData; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaInfoUtil; import org.apache.pulsar.common.util.collections.ConcurrentLongHashMap; @@ -122,7 +121,7 @@ public class ClientCnx extends PulsarHandler { private volatile int numberOfRejectRequests = 0; @Getter - private int maxMessageSize = 0; + private int maxMessageSize = Commands.DEFAULT_MAX_MESSAGE_SIZE; private final int maxNumberOfRejectedRequestPerConnection; private final int rejectedRequestResetTimeSec = 60; @@ -282,9 +281,15 @@ protected void handleConnected(CommandConnected connected) { } checkArgument(state == State.SentConnectFrame || state == State.Connecting); - this.maxMessageSize = connected.getMaxMessageSize(); - ctx.pipeline().replace("defaultFrameDecoder", "frameDecoder", new LengthFieldBasedFrameDecoder( - connected.getMaxMessageSize() + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, 4)); + if (connected.hasMaxMessageSize()) { + if (log.isDebugEnabled()) { + log.debug("{} Connection has max message size setting, replace old frameDecoder with " + + "server frame size {}", ctx.channel(), connected.getMaxMessageSize()); + } + this.maxMessageSize = connected.getMaxMessageSize(); + ctx.pipeline().replace("frameDecoder", "newFrameDecoder", new LengthFieldBasedFrameDecoder( + connected.getMaxMessageSize() + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); + } if (log.isDebugEnabled()) { log.debug("{} Connection is ready", ctx.channel()); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java index 6688bca7d0c18..5d0f87273c50e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarChannelInitializer.java @@ -29,6 +29,7 @@ import org.apache.pulsar.client.api.AuthenticationDataProvider; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.apache.pulsar.common.api.ByteBufPair; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.api.PulsarDecoder; import org.apache.pulsar.common.util.SecurityUtility; @@ -71,7 +72,10 @@ public void initChannel(SocketChannel ch) throws Exception { } ch.pipeline() - .addLast("defaultFrameDecoder", new LengthFieldBasedFrameDecoder(conf.getDefaultMaxFrameSize(), 0, 4, 0, 4)); + .addLast("frameDecoder", + new LengthFieldBasedFrameDecoder( + Commands.DEFAULT_MAX_MESSAGE_SIZE + Commands.MESSAGE_SIZE_FRAME_PADDING, + 0, 4, 0, 4)); ch.pipeline().addLast("handler", clientCnxSupplier.get()); } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java index b57d66b435ead..db77bd17e4ecb 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java @@ -62,7 +62,6 @@ public class ClientConfigurationData implements Serializable, Cloneable { private int connectionTimeoutMs = 10000; private long defaultBackoffIntervalNanos = TimeUnit.MILLISECONDS.toNanos(100); private long maxBackoffIntervalNanos = TimeUnit.SECONDS.toNanos(30); - private int defaultMaxFrameSize = 5242880; public ClientConfigurationData clone() { try { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java index adf8bba1a8ad5..04018f6eed0d1 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java @@ -94,6 +94,11 @@ public class Commands { + // default message size for transfer + public static final int DEFAULT_MAX_MESSAGE_SIZE = 5 * 1024 * 1024; + public static final int MESSAGE_SIZE_FRAME_PADDING = 10 * 1024; + public static final int INVALID_MAX_MESSAGE_SIZE = -1; + public static final short magicCrc32c = 0x0e01; private static final int checksumSize = 4; @@ -189,10 +194,16 @@ public static ByteBuf newConnect(String authMethodName, AuthData authData, int p return res; } + public static ByteBuf newConnected(int clientProtocoVersion) { + return newConnected(clientProtocoVersion, INVALID_MAX_MESSAGE_SIZE); + } + public static ByteBuf newConnected(int clientProtocolVersion, int maxMessageSize) { CommandConnected.Builder connectedBuilder = CommandConnected.newBuilder(); connectedBuilder.setServerVersion("Pulsar Server"); - connectedBuilder.setMaxMessageSize(maxMessageSize); + if (INVALID_MAX_MESSAGE_SIZE != maxMessageSize) { + connectedBuilder.setMaxMessageSize(maxMessageSize); + } // If the broker supports a newer version of the protocol, it will anyway advertise the max version that the // client supports, to avoid confusing the client. diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/PulsarDecoder.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/PulsarDecoder.java index 0e1ea73c45eeb..38ff90f95938f 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/PulsarDecoder.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/PulsarDecoder.java @@ -66,11 +66,6 @@ public abstract class PulsarDecoder extends ChannelInboundHandlerAdapter { - // Max message size is limited by max BookKeeper entry size which is 5MB, and we need to account - // for headers as well. - public final static int MaxMessageSize = (5 * 1024 * 1024 - (10 * 1024)); - public final static int MaxFrameSize = 5 * 1024 * 1024; - @Override public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception { // Get a buffer that contains the full frame diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/raw/MessageParser.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/raw/MessageParser.java index 1f1d66e5bdbca..da24a6c005919 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/raw/MessageParser.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/raw/MessageParser.java @@ -50,7 +50,7 @@ public interface MessageProcessor { * provided {@link MessageProcessor} will be invoked for each individual message. */ public static void parseMessage(TopicName topicName, long ledgerId, long entryId, ByteBuf headersAndPayload, - MessageProcessor processor) throws IOException { + MessageProcessor processor, int maxMessageSize) throws IOException { MessageMetadata msgMetadata = null; ByteBuf payload = headersAndPayload; ByteBuf uncompressedPayload = null; @@ -74,7 +74,7 @@ public static void parseMessage(TopicName topicName, long ledgerId, long entryId } uncompressedPayload = uncompressPayloadIfNeeded(topicName, msgMetadata, headersAndPayload, ledgerId, - entryId); + entryId, maxMessageSize); if (uncompressedPayload == null) { // Message was discarded on decompression error @@ -115,11 +115,11 @@ public static boolean verifyChecksum(TopicName topic, ByteBuf headersAndPayload, } public static ByteBuf uncompressPayloadIfNeeded(TopicName topic, MessageMetadata msgMetadata, - ByteBuf payload, long ledgerId, long entryId) { + ByteBuf payload, long ledgerId, long entryId, int maxMessageSize) { CompressionCodec codec = CompressionCodecProvider.getCompressionCodec(msgMetadata.getCompression()); int uncompressedSize = msgMetadata.getUncompressedSize(); int payloadSize = payload.readableBytes(); - if (payloadSize > PulsarDecoder.MaxMessageSize) { + if (payloadSize > maxMessageSize) { // payload size is itself corrupted since it cannot be bigger than the MaxMessageSize log.error("[{}] Got corrupted payload message size {} at {}:{}", topic, payloadSize, ledgerId, entryId); diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java index a35c6aadedd62..b89a0285b9c86 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServerConnection.java @@ -102,7 +102,7 @@ protected void handleConnect(CommandConnect connect) { return; } } - ctx.writeAndFlush(Commands.newConnected(connect.getProtocolVersion(), ServiceConfig.MAX_MESSAGE_SIZE)); + ctx.writeAndFlush(Commands.newConnected(connect.getProtocolVersion())); state = State.Connected; remoteEndpointProtocolVersion = connect.getProtocolVersion(); } diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java index 0224d91a15613..50d698d076612 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/ServiceChannelInitializer.java @@ -18,7 +18,7 @@ */ package org.apache.pulsar.discovery.service; -import org.apache.pulsar.common.conf.InternalConfigurationData; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.util.NettySslContextBuilder; import org.apache.pulsar.discovery.service.server.ServiceConfig; @@ -37,7 +37,6 @@ public class ServiceChannelInitializer extends ChannelInitializer private final DiscoveryService discoveryService; private final boolean enableTls; private final NettySslContextBuilder sslCtxRefresher; - private final ServiceConfig brokerConf; public ServiceChannelInitializer(DiscoveryService discoveryService, ServiceConfig serviceConfig, boolean e) throws Exception { @@ -53,7 +52,6 @@ public ServiceChannelInitializer(DiscoveryService discoveryService, ServiceConfi } else { this.sslCtxRefresher = null; } - this.brokerConf = serviceConfig; } @Override @@ -65,7 +63,7 @@ protected void initChannel(SocketChannel ch) throws Exception { } } ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( - ServiceConfig.MAX_MESSAGE_SIZE + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, 4)); + Commands.DEFAULT_MAX_MESSAGE_SIZE + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ServerConnection(discoveryService)); } } diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java index 5aedb59735221..94ef92e25d467 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java @@ -33,10 +33,6 @@ */ public class ServiceConfig implements PulsarConfiguration { - // Discovery service doesn't send any messages except Command connected. - // So it's ok use a default value. - public final static int MAX_MESSAGE_SIZE = 5 * 1024 * 1024; - // Local-Zookeeper quorum connection string private String zookeeperServers; // Global-Zookeeper quorum connection string diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index 76d3f7906f7d7..f15d4920787e5 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -98,9 +98,8 @@ protected void initChannel(SocketChannel ch) throws Exception { if (sslCtx != null) { ch.pipeline().addLast(TLS_HANDLER, sslCtx.newHandler(ch.alloc())); } - ch.pipeline().addLast("frameDecoder", - new LengthFieldBasedFrameDecoder(config.getMaxMessagesSize() + 10 * 1024, 0, 4, 0, - 4)); + ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( + Commands.DEFAULT_MAX_MESSAGE_SIZE + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); ch.pipeline().addLast("proxyOutboundHandler", new ProxyBackendHandler(config, protocolVersion)); } }); @@ -260,35 +259,56 @@ protected void handleConnected(CommandConnected connected) { state = BackendState.HandshakeCompleted; - inboundChannel - .writeAndFlush(Commands.newConnected(connected.getProtocolVersion(), connected.getMaxMessageSize())) - .addListener(future -> { - if (log.isDebugEnabled()) { - log.debug("[{}] [{}] Removing decoder from pipeline", inboundChannel, outboundChannel); - } - if (ProxyService.proxyLogLevel == 0) { - // direct tcp proxy - inboundChannel.pipeline().remove("defaultFrameDecoder"); - outboundChannel.pipeline().remove("defaultFrameDecoder"); - } else { - // Enable parsing feature, proxyLogLevel(1 or 2) - // Add parser handler - inboundChannel.pipeline().replace("defaultFrameDecoder", "frameDecoder", - new LengthFieldBasedFrameDecoder( - connected.getMaxMessageSize() - + InternalConfigurationData.MESSAGE_META_SIZE, - 0, 4, 0, 4)); + ChannelFuture channelFuture; + if (connected.hasMaxMessageSize()) { + channelFuture = inboundChannel.writeAndFlush( + Commands.newConnected(connected.getProtocolVersion(), connected.getMaxMessageSize())); + } else { + channelFuture = inboundChannel.writeAndFlush(Commands.newConnected(connected.getProtocolVersion())); + } + + channelFuture.addListener(future -> { + if (log.isDebugEnabled()) { + log.debug("[{}] [{}] Removing decoder from pipeline", inboundChannel, outboundChannel); + } + if (ProxyService.proxyLogLevel == 0) { + // direct tcp proxy + inboundChannel.pipeline().remove("frameDecoder"); + outboundChannel.pipeline().remove("frameDecoder"); + } else { + // Enable parsing feature, proxyLogLevel(1 or 2) + // Add parser handler + if (connected.hasMaxMessageSize()) { + inboundChannel.pipeline().replace("frameDecoder", "newFrameDecoder", + new LengthFieldBasedFrameDecoder(connected.getMaxMessageSize() + + Commands.MESSAGE_SIZE_FRAME_PADDING, + 0, 4, 0, 4)); + outboundChannel.pipeline().replace("frameDecoder", "newFrameDecoder", + new LengthFieldBasedFrameDecoder( + connected.getMaxMessageSize() + + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); + inboundChannel.pipeline().addBefore("handler", "inboundParser", new ParserProxyHandler(inboundChannel, - ParserProxyHandler.FRONTEND_CONN)); + ParserProxyHandler.FRONTEND_CONN, + connected.getMaxMessageSize())); outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", new ParserProxyHandler(outboundChannel, - ParserProxyHandler.BACKEND_CONN)); + ParserProxyHandler.BACKEND_CONN, + connected.getMaxMessageSize())); } - // Start reading from both connections - inboundChannel.read(); - outboundChannel.read(); - }); + inboundChannel.pipeline().addBefore("handler", "inboundParser", + new ParserProxyHandler(inboundChannel, + ParserProxyHandler.FRONTEND_CONN, + Commands.DEFAULT_MAX_MESSAGE_SIZE)); + outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", + new ParserProxyHandler(outboundChannel, + ParserProxyHandler.BACKEND_CONN, + Commands.DEFAULT_MAX_MESSAGE_SIZE)); } + // Start reading from both connections + inboundChannel.read(); + outboundChannel.read(); + }); } @Override diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java index 040a8349b4f88..8b4fe64545347 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java @@ -51,15 +51,18 @@ public class ParserProxyHandler extends ChannelInboundHandlerAdapter { private String connType; + private int maxMessageSize; + //producerid+channelid as key //or consumerid+channelid as key private static Map producerHashMap = new ConcurrentHashMap<>(); private static Map consumerHashMap = new ConcurrentHashMap<>(); - public ParserProxyHandler(Channel channel, String type){ + public ParserProxyHandler(Channel channel, String type, int maxMessageSize){ this.channel = channel; this.connType=type; + this.maxMessageSize = maxMessageSize; } private void logging(Channel conn, PulsarApi.BaseCommand.Type cmdtype, String info, List messages) throws Exception{ @@ -118,7 +121,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) { MessageParser.parseMessage(topicName, -1L, -1L,buffer,(message) -> { messages.add(message); - }); + }, maxMessageSize); logging(ctx.channel() , cmd.getType() , "" , messages); break; @@ -138,7 +141,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) { MessageParser.parseMessage(topicName, -1L, -1L,buffer,(message) -> { messages.add(message); - }); + }, maxMessageSize); logging(ctx.channel() , cmd.getType() , "" , messages); diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java index e728f9045f51e..1ca1159c772c3 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java @@ -111,9 +111,6 @@ public class ProxyConfiguration implements PulsarConfiguration { ) private String brokerWebServiceURLTLS; - @FieldContext(category = CATEGORY_BROKER_DISCOVERY, doc = "") - private int maxMessagesSize = 5233640; - @FieldContext( category = CATEGORY_BROKER_DISCOVERY, doc = "The web service url points to the function worker cluster." diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index a9c86233fd766..18e4e5d750822 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -218,8 +218,7 @@ private void completeConnect() { // partitions metadata lookups state = State.ProxyLookupRequests; lookupProxyHandler = new LookupProxyHandler(service, this); - ctx.writeAndFlush( - Commands.newConnected(protocolVersionToAdvertise, service.getConfiguration().getMaxMessagesSize())); + ctx.writeAndFlush(Commands.newConnected(protocolVersionToAdvertise)); } } diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java index 704518885b0b1..9be109c2027b1 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ServiceChannelInitializer.java @@ -22,7 +22,7 @@ import org.apache.pulsar.client.api.AuthenticationDataProvider; import org.apache.pulsar.client.api.AuthenticationFactory; -import org.apache.pulsar.common.conf.InternalConfigurationData; +import org.apache.pulsar.common.api.Commands; import org.apache.pulsar.common.util.ClientSslContextRefresher; import org.apache.pulsar.common.util.NettySslContextBuilder; @@ -84,9 +84,8 @@ protected void initChannel(SocketChannel ch) throws Exception { } } - ch.pipeline().addLast("defaultFrameDecoder", new LengthFieldBasedFrameDecoder( - proxyService.getConfiguration().getMaxMessagesSize() + InternalConfigurationData.MESSAGE_META_SIZE, 0, 4, 0, - 4)); + ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( + Commands.DEFAULT_MAX_MESSAGE_SIZE + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new ProxyConnection(proxyService, clientSslCtxRefresher == null ? null : clientSslCtxRefresher.get())); } diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java index 23992334ca149..d36c90dbc56fd 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java @@ -23,6 +23,7 @@ import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.bookkeeper.stats.NullStatsProvider; +import org.apache.pulsar.common.api.Commands; import javax.validation.constraints.NotNull; import java.io.IOException; @@ -37,6 +38,7 @@ public class PulsarConnectorConfig implements AutoCloseable { private int targetNumSplits = 2; private int maxSplitMessageQueueSize = 10000; private int maxSplitEntryQueueSize = 1000; + private int maxMessageSize = Commands.DEFAULT_MAX_MESSAGE_SIZE; private String statsProvider = NullStatsProvider.class.getName(); private Map statsProviderConfigs = new HashMap<>(); @@ -59,6 +61,16 @@ public PulsarConnectorConfig setBrokerServiceUrl(String brokerServiceUrl) { return this; } + @Config("pulsar.max-message-size") + public PulsarConnectorConfig setMaxMessageSize(int maxMessageSize) { + this.maxMessageSize = maxMessageSize; + return this; + } + + public int getMaxMessageSize() { + return this.maxMessageSize; + } + @NotNull public String getZookeeperUri() { return this.zookeeperUri; diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java index 4ffbd2f2e80a9..09f3fa624145f 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java @@ -138,6 +138,7 @@ private void initialize(List columnHandles, PulsarSplit puls pulsarSplit.getTableName()); this.metricsTracker = pulsarConnectorMetricsTracker; this.readOffloaded = pulsarConnectorConfig.getManagedLedgerOffloadDriver() != null; + this.pulsarConnectorConfig = pulsarConnectorConfig; Schema schema = PulsarConnectorUtils.parseSchema(pulsarSplit.getSchema()); @@ -260,7 +261,7 @@ public void accept(Entry entry) { } catch (InterruptedException e) { //no-op } - }); + }, pulsarConnectorConfig.getMaxMessageSize()); } catch (IOException e) { log.error(e, "Failed to parse message from pulsar topic %s", topicName.toString()); throw new RuntimeException(e); From 99a2b44b8dbcb8c19d653b79a4ca01f335a95adb Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Mon, 13 May 2019 17:01:17 +0800 Subject: [PATCH 4/8] Fix some error --- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 2 +- .../org/apache/pulsar/proxy/server/DirectProxyHandler.java | 6 ++++-- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 726cba5f031ed..5209e311e24a2 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -495,7 +495,7 @@ public class ServiceConfiguration implements PulsarConfiguration { category = CATEGORY_SERVER, doc = "Max size of messages.", maxValue = Integer.MAX_VALUE - Commands.MESSAGE_SIZE_FRAME_PADDING) - public int maxMessageSize = Commands.DEFAULT_MAX_MESSAGE_SIZE; + private int maxMessageSize = Commands.DEFAULT_MAX_MESSAGE_SIZE; /***** --- TLS --- ****/ @FieldContext( diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index f15d4920787e5..aab3986849c89 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -296,7 +296,7 @@ protected void handleConnected(CommandConnected connected) { new ParserProxyHandler(outboundChannel, ParserProxyHandler.BACKEND_CONN, connected.getMaxMessageSize())); - } + } else { inboundChannel.pipeline().addBefore("handler", "inboundParser", new ParserProxyHandler(inboundChannel, ParserProxyHandler.FRONTEND_CONN, @@ -304,7 +304,9 @@ protected void handleConnected(CommandConnected connected) { outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", new ParserProxyHandler(outboundChannel, ParserProxyHandler.BACKEND_CONN, - Commands.DEFAULT_MAX_MESSAGE_SIZE)); } + Commands.DEFAULT_MAX_MESSAGE_SIZE)); + } + } // Start reading from both connections inboundChannel.read(); outboundChannel.read(); From 59f8d3be61aa47de3a37249c99d04bb09cb3a5ec Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Tue, 14 May 2019 10:32:19 +0800 Subject: [PATCH 5/8] Fix license header --- .../apache/pulsar/discovery/service/server/ServiceConfig.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java index 94ef92e25d467..5f89527a8d241 100644 --- a/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java +++ b/pulsar-discovery-service/src/main/java/org/apache/pulsar/discovery/service/server/ServiceConfig.java @@ -1,4 +1,5 @@ -/** * Licensed to the Apache Software Foundation (ASF) under one +/** + * 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 From 70edbb3b32c5ae29b81919a3ec6b8272d3c8568e Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Wed, 15 May 2019 16:57:45 +0800 Subject: [PATCH 6/8] Add test and make `ClientCnx.maxMessageSize` static --- *Motivation* - Even if the cnx can't use, `maxMessageSize` should be used at compare message size. So it should as a static variable --- .../broker/service/MaxMessageSizeTest.java | 149 ++++++++++++++++++ .../apache/pulsar/client/impl/ClientCnx.java | 4 +- .../pulsar/client/impl/ConsumerImpl.java | 2 +- .../pulsar/client/impl/ProducerImpl.java | 8 +- 4 files changed, 156 insertions(+), 7 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java new file mode 100644 index 0000000000000..e5bc99cf017d1 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java @@ -0,0 +1,149 @@ +/* + * Licensed 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.pulsar.broker.service; + +import com.google.common.collect.Sets; +import java.util.concurrent.TimeUnit; +import org.apache.bookkeeper.conf.ServerConfiguration; +import org.apache.bookkeeper.test.PortManager; +import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.client.admin.PulsarAdmin; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.common.policies.data.ClusterData; +import org.apache.pulsar.common.policies.data.TenantInfo; +import org.apache.pulsar.zookeeper.LocalBookkeeperEnsemble; +import org.testng.Assert; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; + +public class MaxMessageSizeTest { + private static int BROKER_SERVICE_PORT = PortManager.nextFreePort(); + PulsarService pulsar; + ServiceConfiguration configuration; + + PulsarAdmin admin; + + LocalBookkeeperEnsemble bkEnsemble; + + private final int ZOOKEEPER_PORT = PortManager.nextFreePort(); + private final int BROKER_WEBSERVER_PORT = PortManager.nextFreePort(); + + @BeforeMethod + void setup() { + try { + bkEnsemble = new LocalBookkeeperEnsemble(3, ZOOKEEPER_PORT, PortManager::nextFreePort); + ServerConfiguration conf = new ServerConfiguration(); + conf.setNettyMaxFrameSizeBytes(10 * 1024 * 1024); + bkEnsemble.startStandalone(conf, false); + + configuration = new ServiceConfiguration(); + configuration.setZookeeperServers("127.0.0.1:" + ZOOKEEPER_PORT); + configuration.setAdvertisedAddress("localhost"); + configuration.setWebServicePort(BROKER_WEBSERVER_PORT); + configuration.setClusterName("max_message_test"); + configuration.setBrokerServicePort(BROKER_SERVICE_PORT); + configuration.setAuthorizationEnabled(false); + configuration.setAuthenticationEnabled(false); + configuration.setManagedLedgerMaxEntriesPerLedger(5); + configuration.setManagedLedgerMinLedgerRolloverTimeMinutes(0); + configuration.setMaxMessageSize(10 * 1024 * 1024); + + pulsar = new PulsarService(configuration); + pulsar.start(); + + String url = "http://127.0.0.1:" + BROKER_WEBSERVER_PORT; + admin = PulsarAdmin.builder().serviceHttpUrl(url).build(); + admin.clusters().createCluster("max_message_test", new ClusterData(url)); + admin.tenants() + .createTenant("test", new TenantInfo(Sets.newHashSet("appid1"), Sets.newHashSet("max_message_test"))); + admin.namespaces().createNamespace("test/message", Sets.newHashSet("max_message_test")); + } catch (Exception e) { + e.printStackTrace(); + } + } + + @AfterMethod + void shutdown() { + try { + pulsar.close(); + bkEnsemble.stop(); + } catch (Throwable t) { + t.printStackTrace(); + } + } + + @Test + public void testMaxMessageSetting() throws PulsarClientException { + + PulsarClient client = PulsarClient.builder().serviceUrl("pulsar://127.0.0.1:" + BROKER_SERVICE_PORT).build(); + String topicName = "persistent://test/message/topic1"; + Producer producer = client.newProducer().topic(topicName).sendTimeout(60, TimeUnit.SECONDS).create(); + Consumer consumer = client.newConsumer().topic(topicName).subscriptionName("test1").subscribe(); + + // less than 5MB message + + byte[] normalMsg = new byte[2 * 1024 * 1024]; + + try { + producer.send(normalMsg); + } catch (PulsarClientException e) { + Assert.fail("Shouldn't have exception at here", e); + } + + byte[] consumerNormalMsg = consumer.receive().getData(); + Assert.assertEquals(normalMsg, consumerNormalMsg); + + // equal 5MB message + byte[] limitMsg = new byte[5 * 1024 * 1024]; + try { + producer.send(limitMsg); + } catch (PulsarClientException e) { + Assert.fail("Shouldn't have exception at here", e); + } + + byte[] consumerLimitMsg = consumer.receive().getData(); + Assert.assertEquals(limitMsg, consumerLimitMsg); + + // less than 10MB message + byte[] newNormalMsg = new byte[8 * 1024 * 1024]; + try { + producer.send(newNormalMsg); + } catch (PulsarClientException e) { + Assert.fail("Shouldn't have exception at here", e); + } + + byte[] consumerNewNormalMsg = consumer.receive().getData(); + Assert.assertEquals(newNormalMsg, consumerNewNormalMsg); + + // equals 10MB message + byte[] newLimitMsg = new byte[10 * 1024 * 1024]; + try { + producer.send(newLimitMsg); + Assert.fail("Shouldn't send out this message"); + } catch (PulsarClientException e) { + // no-op + } + + consumer.unsubscribe(); + consumer.close(); + producer.close(); + client.close(); + + } +} diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index b3ef8b2be35cd..7a58357ef9720 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -121,7 +121,7 @@ public class ClientCnx extends PulsarHandler { private volatile int numberOfRejectRequests = 0; @Getter - private int maxMessageSize = Commands.DEFAULT_MAX_MESSAGE_SIZE; + private static int maxMessageSize = Commands.DEFAULT_MAX_MESSAGE_SIZE; private final int maxNumberOfRejectedRequestPerConnection; private final int rejectedRequestResetTimeSec = 60; @@ -286,7 +286,7 @@ protected void handleConnected(CommandConnected connected) { log.debug("{} Connection has max message size setting, replace old frameDecoder with " + "server frame size {}", ctx.channel(), connected.getMaxMessageSize()); } - this.maxMessageSize = connected.getMaxMessageSize(); + maxMessageSize = connected.getMaxMessageSize(); ctx.pipeline().replace("frameDecoder", "newFrameDecoder", new LengthFieldBasedFrameDecoder( connected.getMaxMessageSize() + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 9946c48ee3122..b2421f246335c 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -1140,7 +1140,7 @@ private ByteBuf uncompressPayloadIfNeeded(MessageIdData messageId, MessageMetada CompressionCodec codec = CompressionCodecProvider.getCompressionCodec(compressionType); int uncompressedSize = msgMetadata.getUncompressedSize(); int payloadSize = payload.readableBytes(); - if (payloadSize > cnx().getMaxMessageSize()) { + if (payloadSize > ClientCnx.getMaxMessageSize()) { // payload size is itself corrupted since it cannot be bigger than the MaxMessageSize log.error("[{}][{}] Got corrupted payload message size {} at {}", topic, subscription, payloadSize, messageId); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 831eba2f7b2a3..3e32fb91d2e3e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -311,14 +311,14 @@ public void sendAsync(Message message, SendCallback callback) { // validate msg-size (For batching this will be check at the batch completion size) int compressedSize = compressedPayload.readableBytes(); - if (compressedSize > cnx().getMaxMessageSize()) { + if (compressedSize > ClientCnx.getMaxMessageSize()) { compressedPayload.release(); String compressedStr = (!isBatchMessagingEnabled() && conf.getCompressionType() != CompressionType.NONE) ? "Compressed" : ""; PulsarClientException.InvalidMessageException invalidMessageException = new PulsarClientException.InvalidMessageException( format("%s Message payload size %d cannot exceed %d bytes", compressedStr, compressedSize, - cnx().getMaxMessageSize())); + ClientCnx.getMaxMessageSize())); callback.sendComplete(invalidMessageException); return; } @@ -1304,12 +1304,12 @@ private void batchMessageAndSend() { op = OpSendMsg.create(batchMessageContainer.messages, cmd, sequenceId, batchMessageContainer.firstCallback); - if (encryptedPayload.readableBytes() > cnx().getMaxMessageSize()) { + if (encryptedPayload.readableBytes() > ClientCnx.getMaxMessageSize()) { cmd.release(); semaphore.release(numMessagesInBatch); if (op != null) { op.callback.sendComplete(new PulsarClientException.InvalidMessageException( - "Message size is bigger than " + cnx().getMaxMessageSize() + " bytes")); + "Message size is bigger than " + ClientCnx.getMaxMessageSize() + " bytes")); } return; } From 8847c3284468361af2531be8e501a3c7059e7584 Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Wed, 15 May 2019 17:04:05 +0800 Subject: [PATCH 7/8] fix code style --- .../common/conf/InternalConfigurationData.java | 1 - .../pulsar/proxy/server/DirectProxyHandler.java | 16 ++++++++-------- 2 files changed, 8 insertions(+), 9 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java b/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java index 45d9c110796c4..558c1aa8da8a5 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java @@ -23,7 +23,6 @@ public class InternalConfigurationData { - public final static int MESSAGE_META_SIZE = 10 * 1024; private String zookeeperServers; private String configurationStoreServers; private String ledgersRootPath; diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index aab3986849c89..a8deb4425e175 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -297,14 +297,14 @@ protected void handleConnected(CommandConnected connected) { ParserProxyHandler.BACKEND_CONN, connected.getMaxMessageSize())); } else { - inboundChannel.pipeline().addBefore("handler", "inboundParser", - new ParserProxyHandler(inboundChannel, - ParserProxyHandler.FRONTEND_CONN, - Commands.DEFAULT_MAX_MESSAGE_SIZE)); - outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", - new ParserProxyHandler(outboundChannel, - ParserProxyHandler.BACKEND_CONN, - Commands.DEFAULT_MAX_MESSAGE_SIZE)); + inboundChannel.pipeline().addBefore("handler", "inboundParser", + new ParserProxyHandler(inboundChannel, + ParserProxyHandler.FRONTEND_CONN, + Commands.DEFAULT_MAX_MESSAGE_SIZE)); + outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", + new ParserProxyHandler(outboundChannel, + ParserProxyHandler.BACKEND_CONN, + Commands.DEFAULT_MAX_MESSAGE_SIZE)); } } // Start reading from both connections From fb83d3fcb23dab180d7b583f71daca091f26399e Mon Sep 17 00:00:00 2001 From: Yong Zhang Date: Wed, 15 May 2019 17:44:53 +0800 Subject: [PATCH 8/8] Fix license header --- .../broker/service/MaxMessageSizeTest.java | 25 +++++++++++-------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java index e5bc99cf017d1..4fb3183ab5a39 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MaxMessageSizeTest.java @@ -1,15 +1,20 @@ -/* - * Licensed 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 +/** + * 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 + * 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. + * 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.pulsar.broker.service;