mirror of https://github.com/apache/cassandra
Update Netty dependencies to latest, clean up SocketFactory
patch by Aleksey Yeschenko; reviewed by Benedict Elliott Smith for CASSANDRA-15195
This commit is contained in:
parent
0b8c5f70a2
commit
a339aa9e98
|
|
@ -1,4 +1,5 @@
|
|||
4.0
|
||||
* Update Netty dependencies to latest, clean up SocketFactory (CASSANDRA-15195)
|
||||
* Native Transport - Apply noSpamLogger to ConnectionLimitHandler (CASSANDRA-15167)
|
||||
* Reduce heap pressure during compactions (CASSANDRA-14654)
|
||||
* Support building Cassandra with JDK 11 (CASSANDRA-15108)
|
||||
|
|
|
|||
|
|
@ -548,7 +548,8 @@
|
|||
<dependency groupId="com.addthis.metrics" artifactId="reporter-config3" version="3.0.3" />
|
||||
<dependency groupId="org.mindrot" artifactId="jbcrypt" version="0.3m" />
|
||||
<dependency groupId="io.airlift" artifactId="airline" version="0.8" />
|
||||
<dependency groupId="io.netty" artifactId="netty-all" version="4.1.28.Final" />
|
||||
<dependency groupId="io.netty" artifactId="netty-all" version="4.1.37.Final" />
|
||||
<dependency groupId="io.netty" artifactId="netty-tcnative-boringssl-static" version="2.0.25.Final" />
|
||||
<dependency groupId="net.openhft" artifactId="chronicle-queue" version="${chronicle-queue.version}"/>
|
||||
<dependency groupId="net.openhft" artifactId="chronicle-core" version="${chronicle-core.version}"/>
|
||||
<dependency groupId="net.openhft" artifactId="chronicle-bytes" version="${chronicle-bytes.version}"/>
|
||||
|
|
|
|||
|
|
@ -0,0 +1,201 @@
|
|||
Apache License
|
||||
Version 2.0, January 2004
|
||||
http://www.apache.org/licenses/
|
||||
|
||||
TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
|
||||
|
||||
1. Definitions.
|
||||
|
||||
"License" shall mean the terms and conditions for use, reproduction,
|
||||
and distribution as defined by Sections 1 through 9 of this document.
|
||||
|
||||
"Licensor" shall mean the copyright owner or entity authorized by
|
||||
the copyright owner that is granting the License.
|
||||
|
||||
"Legal Entity" shall mean the union of the acting entity and all
|
||||
other entities that control, are controlled by, or are under common
|
||||
control with that entity. For the purposes of this definition,
|
||||
"control" means (i) the power, direct or indirect, to cause the
|
||||
direction or management of such entity, whether by contract or
|
||||
otherwise, or (ii) ownership of fifty percent (50%) or more of the
|
||||
outstanding shares, or (iii) beneficial ownership of such entity.
|
||||
|
||||
"You" (or "Your") shall mean an individual or Legal Entity
|
||||
exercising permissions granted by this License.
|
||||
|
||||
"Source" form shall mean the preferred form for making modifications,
|
||||
including but not limited to software source code, documentation
|
||||
source, and configuration files.
|
||||
|
||||
"Object" form shall mean any form resulting from mechanical
|
||||
transformation or translation of a Source form, including but
|
||||
not limited to compiled object code, generated documentation,
|
||||
and conversions to other media types.
|
||||
|
||||
"Work" shall mean the work of authorship, whether in Source or
|
||||
Object form, made available under the License, as indicated by a
|
||||
copyright notice that is included in or attached to the work
|
||||
(an example is provided in the Appendix below).
|
||||
|
||||
"Derivative Works" shall mean any work, whether in Source or Object
|
||||
form, that is based on (or derived from) the Work and for which the
|
||||
editorial revisions, annotations, elaborations, or other modifications
|
||||
represent, as a whole, an original work of authorship. For the purposes
|
||||
of this License, Derivative Works shall not include works that remain
|
||||
separable from, or merely link (or bind by name) to the interfaces of,
|
||||
the Work and Derivative Works thereof.
|
||||
|
||||
"Contribution" shall mean any work of authorship, including
|
||||
the original version of the Work and any modifications or additions
|
||||
to that Work or Derivative Works thereof, that is intentionally
|
||||
submitted to Licensor for inclusion in the Work by the copyright owner
|
||||
or by an individual or Legal Entity authorized to submit on behalf of
|
||||
the copyright owner. For the purposes of this definition, "submitted"
|
||||
means any form of electronic, verbal, or written communication sent
|
||||
to the Licensor or its representatives, including but not limited to
|
||||
communication on electronic mailing lists, source code control systems,
|
||||
and issue tracking systems that are managed by, or on behalf of, the
|
||||
Licensor for the purpose of discussing and improving the Work, but
|
||||
excluding communication that is conspicuously marked or otherwise
|
||||
designated in writing by the copyright owner as "Not a Contribution."
|
||||
|
||||
"Contributor" shall mean Licensor and any individual or Legal Entity
|
||||
on behalf of whom a Contribution has been received by Licensor and
|
||||
subsequently incorporated within the Work.
|
||||
|
||||
2. Grant of Copyright License. Subject to the terms and conditions of
|
||||
this License, each Contributor hereby grants to You a perpetual,
|
||||
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
|
||||
copyright license to reproduce, prepare Derivative Works of,
|
||||
publicly display, publicly perform, sublicense, and distribute the
|
||||
Work and such Derivative Works in Source or Object form.
|
||||
|
||||
3. Grant of Patent License. Subject to the terms and conditions of
|
||||
this License, each Contributor hereby grants to You a perpetual,
|
||||
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
|
||||
(except as stated in this section) patent license to make, have made,
|
||||
use, offer to sell, sell, import, and otherwise transfer the Work,
|
||||
where such license applies only to those patent claims licensable
|
||||
by such Contributor that are necessarily infringed by their
|
||||
Contribution(s) alone or by combination of their Contribution(s)
|
||||
with the Work to which such Contribution(s) was submitted. If You
|
||||
institute patent litigation against any entity (including a
|
||||
cross-claim or counterclaim in a lawsuit) alleging that the Work
|
||||
or a Contribution incorporated within the Work constitutes direct
|
||||
or contributory patent infringement, then any patent licenses
|
||||
granted to You under this License for that Work shall terminate
|
||||
as of the date such litigation is filed.
|
||||
|
||||
4. Redistribution. You may reproduce and distribute copies of the
|
||||
Work or Derivative Works thereof in any medium, with or without
|
||||
modifications, and in Source or Object form, provided that You
|
||||
meet the following conditions:
|
||||
|
||||
(a) You must give any other recipients of the Work or
|
||||
Derivative Works a copy of this License; and
|
||||
|
||||
(b) You must cause any modified files to carry prominent notices
|
||||
stating that You changed the files; and
|
||||
|
||||
(c) You must retain, in the Source form of any Derivative Works
|
||||
that You distribute, all copyright, patent, trademark, and
|
||||
attribution notices from the Source form of the Work,
|
||||
excluding those notices that do not pertain to any part of
|
||||
the Derivative Works; and
|
||||
|
||||
(d) If the Work includes a "NOTICE" text file as part of its
|
||||
distribution, then any Derivative Works that You distribute must
|
||||
include a readable copy of the attribution notices contained
|
||||
within such NOTICE file, excluding those notices that do not
|
||||
pertain to any part of the Derivative Works, in at least one
|
||||
of the following places: within a NOTICE text file distributed
|
||||
as part of the Derivative Works; within the Source form or
|
||||
documentation, if provided along with the Derivative Works; or,
|
||||
within a display generated by the Derivative Works, if and
|
||||
wherever such third-party notices normally appear. The contents
|
||||
of the NOTICE file are for informational purposes only and
|
||||
do not modify the License. You may add Your own attribution
|
||||
notices within Derivative Works that You distribute, alongside
|
||||
or as an addendum to the NOTICE text from the Work, provided
|
||||
that such additional attribution notices cannot be construed
|
||||
as modifying the License.
|
||||
|
||||
You may add Your own copyright statement to Your modifications and
|
||||
may provide additional or different license terms and conditions
|
||||
for use, reproduction, or distribution of Your modifications, or
|
||||
for any such Derivative Works as a whole, provided Your use,
|
||||
reproduction, and distribution of the Work otherwise complies with
|
||||
the conditions stated in this License.
|
||||
|
||||
5. Submission of Contributions. Unless You explicitly state otherwise,
|
||||
any Contribution intentionally submitted for inclusion in the Work
|
||||
by You to the Licensor shall be under the terms and conditions of
|
||||
this License, without any additional terms or conditions.
|
||||
Notwithstanding the above, nothing herein shall supersede or modify
|
||||
the terms of any separate license agreement you may have executed
|
||||
with Licensor regarding such Contributions.
|
||||
|
||||
6. Trademarks. This License does not grant permission to use the trade
|
||||
names, trademarks, service marks, or product names of the Licensor,
|
||||
except as required for reasonable and customary use in describing the
|
||||
origin of the Work and reproducing the content of the NOTICE file.
|
||||
|
||||
7. Disclaimer of Warranty. Unless required by applicable law or
|
||||
agreed to in writing, Licensor provides the Work (and each
|
||||
Contributor provides its Contributions) on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
|
||||
implied, including, without limitation, any warranties or conditions
|
||||
of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
|
||||
PARTICULAR PURPOSE. You are solely responsible for determining the
|
||||
appropriateness of using or redistributing the Work and assume any
|
||||
risks associated with Your exercise of permissions under this License.
|
||||
|
||||
8. Limitation of Liability. In no event and under no legal theory,
|
||||
whether in tort (including negligence), contract, or otherwise,
|
||||
unless required by applicable law (such as deliberate and grossly
|
||||
negligent acts) or agreed to in writing, shall any Contributor be
|
||||
liable to You for damages, including any direct, indirect, special,
|
||||
incidental, or consequential damages of any character arising as a
|
||||
result of this License or out of the use or inability to use the
|
||||
Work (including but not limited to damages for loss of goodwill,
|
||||
work stoppage, computer failure or malfunction, or any and all
|
||||
other commercial damages or losses), even if such Contributor
|
||||
has been advised of the possibility of such damages.
|
||||
|
||||
9. Accepting Warranty or Additional Liability. While redistributing
|
||||
the Work or Derivative Works thereof, You may choose to offer,
|
||||
and charge a fee for, acceptance of support, warranty, indemnity,
|
||||
or other liability obligations and/or rights consistent with this
|
||||
License. However, in accepting such obligations, You may act only
|
||||
on Your own behalf and on Your sole responsibility, not on behalf
|
||||
of any other Contributor, and only if You agree to indemnify,
|
||||
defend, and hold each Contributor harmless for any liability
|
||||
incurred by, or claims asserted against, such Contributor by reason
|
||||
of your accepting any such warranty or additional liability.
|
||||
|
||||
END OF TERMS AND CONDITIONS
|
||||
|
||||
APPENDIX: How to apply the Apache License to your work.
|
||||
|
||||
To apply the Apache License to your work, attach the following
|
||||
boilerplate notice, with the fields enclosed by brackets "[]"
|
||||
replaced with your own identifying information. (Don't include
|
||||
the brackets!) The text should be enclosed in the appropriate
|
||||
comment syntax for the file format. We also recommend that a
|
||||
file or class name and description of purpose be included on the
|
||||
same "printed page" as the copyright notice for easier
|
||||
identification within third-party archives.
|
||||
|
||||
Copyright [yyyy] [name of copyright owner]
|
||||
|
||||
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.
|
||||
Binary file not shown.
Binary file not shown.
|
|
@ -132,6 +132,8 @@ public class InboundConnectionInitiator
|
|||
ServerBootstrap bootstrap = initializer.settings.socketFactory
|
||||
.newServerBootstrap()
|
||||
.option(ChannelOption.SO_BACKLOG, 1 << 9)
|
||||
.option(ChannelOption.ALLOCATOR, GlobalBufferPoolAllocator.instance)
|
||||
.option(ChannelOption.SO_REUSEADDR, true)
|
||||
.childHandler(initializer);
|
||||
|
||||
int socketReceiveBufferSizeInBytes = initializer.settings.socketReceiveBufferSizeInBytes;
|
||||
|
|
|
|||
|
|
@ -171,13 +171,15 @@ public class OutboundConnectionInitiator<SuccessType extends OutboundConnectionI
|
|||
*/
|
||||
private Bootstrap createBootstrap(EventLoop eventLoop)
|
||||
{
|
||||
Bootstrap bootstrap = newBootstrap(eventLoop, settings.tcpUserTimeoutInMS)
|
||||
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, settings.tcpConnectTimeoutInMS)
|
||||
.option(ChannelOption.SO_KEEPALIVE, true)
|
||||
.option(ChannelOption.SO_REUSEADDR, true)
|
||||
.option(ChannelOption.TCP_NODELAY, settings.tcpNoDelay)
|
||||
.option(ChannelOption.MESSAGE_SIZE_ESTIMATOR, NoSizeEstimator.instance)
|
||||
.handler(new Initializer());
|
||||
Bootstrap bootstrap = settings.socketFactory
|
||||
.newClientBootstrap(eventLoop, settings.tcpUserTimeoutInMS)
|
||||
.option(ChannelOption.ALLOCATOR, GlobalBufferPoolAllocator.instance)
|
||||
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, settings.tcpConnectTimeoutInMS)
|
||||
.option(ChannelOption.SO_KEEPALIVE, true)
|
||||
.option(ChannelOption.SO_REUSEADDR, true)
|
||||
.option(ChannelOption.TCP_NODELAY, settings.tcpNoDelay)
|
||||
.option(ChannelOption.MESSAGE_SIZE_ESTIMATOR, NoSizeEstimator.instance)
|
||||
.handler(new Initializer());
|
||||
|
||||
if (settings.socketSendBufferSizeInBytes > 0)
|
||||
bootstrap.option(ChannelOption.SO_SNDBUF, settings.socketSendBufferSizeInBytes);
|
||||
|
|
|
|||
|
|
@ -18,13 +18,14 @@
|
|||
package org.apache.cassandra.net;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.Field;
|
||||
import java.net.ConnectException;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.nio.channels.ClosedChannelException;
|
||||
import java.nio.channels.spi.SelectorProvider;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ThreadFactory;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import javax.annotation.Nullable;
|
||||
import javax.net.ssl.SSLEngine;
|
||||
|
|
@ -37,11 +38,11 @@ import org.slf4j.LoggerFactory;
|
|||
import io.netty.bootstrap.Bootstrap;
|
||||
import io.netty.bootstrap.ServerBootstrap;
|
||||
import io.netty.channel.Channel;
|
||||
import io.netty.channel.ChannelOption;
|
||||
import io.netty.channel.ChannelFactory;
|
||||
import io.netty.channel.DefaultSelectStrategyFactory;
|
||||
import io.netty.channel.EventLoop;
|
||||
import io.netty.channel.EventLoopGroup;
|
||||
import io.netty.channel.MultithreadEventLoopGroup;
|
||||
import io.netty.channel.SingleThreadEventLoop;
|
||||
import io.netty.channel.ServerChannel;
|
||||
import io.netty.channel.epoll.EpollChannelOption;
|
||||
import io.netty.channel.epoll.EpollEventLoopGroup;
|
||||
import io.netty.channel.epoll.EpollServerSocketChannel;
|
||||
|
|
@ -53,10 +54,10 @@ import io.netty.channel.unix.Errors;
|
|||
import io.netty.handler.ssl.OpenSsl;
|
||||
import io.netty.handler.ssl.SslContext;
|
||||
import io.netty.handler.ssl.SslHandler;
|
||||
import io.netty.util.concurrent.DefaultEventExecutorChooserFactory;
|
||||
import io.netty.util.concurrent.DefaultThreadFactory;
|
||||
import io.netty.util.concurrent.EventExecutor;
|
||||
import io.netty.util.concurrent.MultithreadEventExecutorGroup;
|
||||
import io.netty.util.concurrent.SingleThreadEventExecutor;
|
||||
import io.netty.util.concurrent.RejectedExecutionHandlers;
|
||||
import io.netty.util.concurrent.ThreadPerTaskExecutor;
|
||||
import io.netty.util.internal.logging.InternalLoggerFactory;
|
||||
import io.netty.util.internal.logging.Slf4JLoggerFactory;
|
||||
import org.apache.cassandra.concurrent.NamedThreadFactory;
|
||||
|
|
@ -82,8 +83,85 @@ public final class SocketFactory
|
|||
|
||||
private static final int EVENT_THREADS = Integer.getInteger(Config.PROPERTY_PREFIX + "internode-event-threads", FBUtilities.getAvailableProcessors());
|
||||
|
||||
public enum Provider { EPOLL, NIO }
|
||||
private static final Provider DEFAULT_PROVIDER = NativeTransportService.useEpoll() ? Provider.EPOLL : Provider.NIO;
|
||||
/**
|
||||
* The default task queue used by {@code NioEventLoop} and {@code EpollEventLoop} is {@code MpscUnboundedArrayQueue},
|
||||
* provided by JCTools. While efficient, it has an undesirable quality for a queue backing an event loop: it is
|
||||
* not non-blocking, and can cause the event loop to busy-spin while waiting for a partially completed task
|
||||
* offer, if the producer thread has been suspended mid-offer.
|
||||
*
|
||||
* As it happens, however, we have an MPSC queue implementation that is perfectly fit for this purpose -
|
||||
* {@link ManyToOneConcurrentLinkedQueue}, that is non-blocking, and already used throughout the codebase,
|
||||
* that we can and do use here as well.
|
||||
*/
|
||||
enum Provider
|
||||
{
|
||||
NIO
|
||||
{
|
||||
@Override
|
||||
NioEventLoopGroup makeEventLoopGroup(int threadCount, ThreadFactory threadFactory)
|
||||
{
|
||||
return new NioEventLoopGroup(threadCount,
|
||||
new ThreadPerTaskExecutor(threadFactory),
|
||||
DefaultEventExecutorChooserFactory.INSTANCE,
|
||||
SelectorProvider.provider(),
|
||||
DefaultSelectStrategyFactory.INSTANCE,
|
||||
RejectedExecutionHandlers.reject(),
|
||||
capacity -> new ManyToOneConcurrentLinkedQueue<>());
|
||||
}
|
||||
|
||||
@Override
|
||||
ChannelFactory<NioSocketChannel> clientChannelFactory()
|
||||
{
|
||||
return NioSocketChannel::new;
|
||||
}
|
||||
|
||||
@Override
|
||||
ChannelFactory<NioServerSocketChannel> serverChannelFactory()
|
||||
{
|
||||
return NioServerSocketChannel::new;
|
||||
}
|
||||
},
|
||||
EPOLL
|
||||
{
|
||||
@Override
|
||||
EpollEventLoopGroup makeEventLoopGroup(int threadCount, ThreadFactory threadFactory)
|
||||
{
|
||||
return new EpollEventLoopGroup(threadCount,
|
||||
new ThreadPerTaskExecutor(threadFactory),
|
||||
DefaultEventExecutorChooserFactory.INSTANCE,
|
||||
DefaultSelectStrategyFactory.INSTANCE,
|
||||
RejectedExecutionHandlers.reject(),
|
||||
capacity -> new ManyToOneConcurrentLinkedQueue<>());
|
||||
}
|
||||
|
||||
@Override
|
||||
ChannelFactory<EpollSocketChannel> clientChannelFactory()
|
||||
{
|
||||
return EpollSocketChannel::new;
|
||||
}
|
||||
|
||||
@Override
|
||||
ChannelFactory<EpollServerSocketChannel> serverChannelFactory()
|
||||
{
|
||||
return EpollServerSocketChannel::new;
|
||||
}
|
||||
};
|
||||
|
||||
EventLoopGroup makeEventLoopGroup(int threadCount, String threadNamePrefix)
|
||||
{
|
||||
logger.debug("using netty {} event loop for pool prefix {}", name(), threadNamePrefix);
|
||||
return makeEventLoopGroup(threadCount, new DefaultThreadFactory(threadNamePrefix, true));
|
||||
}
|
||||
|
||||
abstract EventLoopGroup makeEventLoopGroup(int threadCount, ThreadFactory threadFactory);
|
||||
abstract ChannelFactory<? extends Channel> clientChannelFactory();
|
||||
abstract ChannelFactory<? extends ServerChannel> serverChannelFactory();
|
||||
|
||||
static Provider optimalProvider()
|
||||
{
|
||||
return NativeTransportService.useEpoll() ? EPOLL : NIO;
|
||||
}
|
||||
}
|
||||
|
||||
/** a useful addition for debugging; simply set to true to get more data in your logs */
|
||||
static final boolean WIRETRACE = false;
|
||||
|
|
@ -93,94 +171,42 @@ public final class SocketFactory
|
|||
InternalLoggerFactory.setDefaultFactory(Slf4JLoggerFactory.INSTANCE);
|
||||
}
|
||||
|
||||
private final Provider provider;
|
||||
private final EventLoopGroup acceptGroup;
|
||||
private final EventLoopGroup defaultGroup;
|
||||
// we need a separate EventLoopGroup for outbound streaming because sendFile is blocking
|
||||
private final EventLoopGroup outboundStreamingGroup;
|
||||
final ExecutorService synchronousWorkExecutor = Executors.newCachedThreadPool(new NamedThreadFactory("Messaging-SynchronousWork"));
|
||||
|
||||
SocketFactory() { this(DEFAULT_PROVIDER); }
|
||||
SocketFactory()
|
||||
{
|
||||
this(Provider.optimalProvider());
|
||||
}
|
||||
|
||||
SocketFactory(Provider provider)
|
||||
{
|
||||
this.acceptGroup = getEventLoopGroup(provider, 1, "Messaging-AcceptLoop");
|
||||
this.defaultGroup = getEventLoopGroup(provider, EVENT_THREADS, NamedThreadFactory.globalPrefix() + "Messaging-EventLoop");
|
||||
this.outboundStreamingGroup = getEventLoopGroup(provider, EVENT_THREADS, "Streaming-EventLoop");
|
||||
assert provider == providerOf(acceptGroup)
|
||||
&& provider == providerOf(defaultGroup)
|
||||
&& provider == providerOf(outboundStreamingGroup);
|
||||
this.provider = provider;
|
||||
this.acceptGroup = provider.makeEventLoopGroup(1, "Messaging-AcceptLoop");
|
||||
this.defaultGroup = provider.makeEventLoopGroup(EVENT_THREADS, NamedThreadFactory.globalPrefix() + "Messaging-EventLoop");
|
||||
this.outboundStreamingGroup = provider.makeEventLoopGroup(EVENT_THREADS, "Streaming-EventLoop");
|
||||
}
|
||||
|
||||
private static EventLoopGroup getEventLoopGroup(Provider provider, int threadCount, String threadNamePrefix)
|
||||
{
|
||||
switch (provider)
|
||||
{
|
||||
case EPOLL:
|
||||
logger.debug("using netty epoll event loop for pool prefix {}", threadNamePrefix);
|
||||
return overwriteMPSCQueues(new EpollEventLoopGroup(threadCount, new DefaultThreadFactory(threadNamePrefix, true)));
|
||||
case NIO:
|
||||
logger.debug("using netty nio event loop for pool prefix {}", threadNamePrefix);
|
||||
return overwriteMPSCQueues(new NioEventLoopGroup(threadCount, new DefaultThreadFactory(threadNamePrefix, true)));
|
||||
default:
|
||||
throw new IllegalStateException();
|
||||
}
|
||||
}
|
||||
|
||||
private static Provider providerOf(EventLoopGroup eventLoopGroup)
|
||||
{
|
||||
while (eventLoopGroup instanceof SingleThreadEventLoop)
|
||||
eventLoopGroup = ((SingleThreadEventLoop) eventLoopGroup).parent();
|
||||
|
||||
if (eventLoopGroup instanceof EpollEventLoopGroup)
|
||||
return Provider.EPOLL;
|
||||
if (eventLoopGroup instanceof NioEventLoopGroup)
|
||||
return Provider.NIO;
|
||||
throw new IllegalStateException();
|
||||
}
|
||||
|
||||
static Bootstrap newBootstrap(EventLoop eventLoop, int tcpUserTimeoutInMS)
|
||||
Bootstrap newClientBootstrap(EventLoop eventLoop, int tcpUserTimeoutInMS)
|
||||
{
|
||||
if (eventLoop == null)
|
||||
throw new IllegalArgumentException("must provide eventLoop");
|
||||
|
||||
Bootstrap bootstrap = new Bootstrap()
|
||||
.group(eventLoop)
|
||||
.option(ChannelOption.ALLOCATOR, GlobalBufferPoolAllocator.instance)
|
||||
.option(ChannelOption.SO_KEEPALIVE, true);
|
||||
Bootstrap bootstrap = new Bootstrap().group(eventLoop).channelFactory(provider.clientChannelFactory());
|
||||
|
||||
if (provider == Provider.EPOLL)
|
||||
bootstrap.option(EpollChannelOption.TCP_USER_TIMEOUT, tcpUserTimeoutInMS);
|
||||
|
||||
switch (providerOf(eventLoop))
|
||||
{
|
||||
case EPOLL:
|
||||
bootstrap.channel(EpollSocketChannel.class);
|
||||
bootstrap.option(EpollChannelOption.TCP_USER_TIMEOUT, tcpUserTimeoutInMS);
|
||||
break;
|
||||
case NIO:
|
||||
bootstrap.channel(NioSocketChannel.class);
|
||||
}
|
||||
return bootstrap;
|
||||
}
|
||||
|
||||
ServerBootstrap newServerBootstrap()
|
||||
{
|
||||
return newServerBootstrap(acceptGroup, defaultGroup);
|
||||
}
|
||||
|
||||
private static ServerBootstrap newServerBootstrap(EventLoopGroup acceptGroup, EventLoopGroup defaultGroup)
|
||||
{
|
||||
ServerBootstrap bootstrap = new ServerBootstrap()
|
||||
.group(acceptGroup, defaultGroup)
|
||||
.option(ChannelOption.ALLOCATOR, GlobalBufferPoolAllocator.instance)
|
||||
.option(ChannelOption.SO_REUSEADDR, true);
|
||||
|
||||
switch (providerOf(defaultGroup))
|
||||
{
|
||||
case EPOLL:
|
||||
bootstrap.channel(EpollServerSocketChannel.class);
|
||||
break;
|
||||
case NIO:
|
||||
bootstrap.channel(NioServerSocketChannel.class);
|
||||
}
|
||||
|
||||
return bootstrap;
|
||||
return new ServerBootstrap().group(acceptGroup, defaultGroup).channelFactory(provider.serverChannelFactory());
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -270,58 +296,4 @@ public final class SocketFactory
|
|||
{
|
||||
return from + "->" + to + '-' + type + '-' + id;
|
||||
}
|
||||
|
||||
/**
|
||||
* The default task queue used by {@code NioEventLoop} and {@code EpollEventLoop} is {@code MpscUnboundedArrayQueue},
|
||||
* provided by JCTools. While efficient, it has an undesirable quality for a queue backing an event loop: it is
|
||||
* not non-blocking, and can cause the event loop to busy-spin while waiting for a partially completed task
|
||||
* offer, if the producer thread has been suspended mid-offer. Sadly, there is currently no way to work around
|
||||
* this behaviour in application-logic.
|
||||
*
|
||||
* As it happens, however, we have an MPSC queue implementation that is perfectly fit for this purpose -
|
||||
* {@link ManyToOneConcurrentLinkedQueue}, that is non-blocking, and already used throughout the codebase.
|
||||
*
|
||||
* Unfortunately, there is no Netty API or to override the default queue, so we have to resort to reflection,
|
||||
* for now.
|
||||
*
|
||||
* We filed a Netty issue asking for this capability to be provided cleanly:
|
||||
* https://github.com/netty/netty/issues/9105, and hopefully Netty will implement it some day. When and if
|
||||
* that happens, this reflection-based workaround should be removed.
|
||||
*/
|
||||
private static EventLoopGroup overwriteMPSCQueues(MultithreadEventLoopGroup eventLoopGroup)
|
||||
{
|
||||
try
|
||||
{
|
||||
for (EventExecutor eventExecutor : (EventExecutor[]) childrenField.get(eventLoopGroup))
|
||||
{
|
||||
SingleThreadEventLoop eventLoop = (SingleThreadEventLoop) eventExecutor;
|
||||
taskQueueField.set(eventLoop, new ManyToOneConcurrentLinkedQueue<>());
|
||||
tailTasksField.set(eventLoop, new ManyToOneConcurrentLinkedQueue<>());
|
||||
}
|
||||
return eventLoopGroup;
|
||||
}
|
||||
catch (IllegalAccessException e)
|
||||
{
|
||||
throw new IllegalStateException(e);
|
||||
}
|
||||
}
|
||||
|
||||
private static final Field childrenField, taskQueueField, tailTasksField;
|
||||
static
|
||||
{
|
||||
try
|
||||
{
|
||||
childrenField = MultithreadEventExecutorGroup.class.getDeclaredField("children");
|
||||
taskQueueField = SingleThreadEventExecutor.class.getDeclaredField("taskQueue");
|
||||
tailTasksField = SingleThreadEventLoop.class.getDeclaredField("tailTasks");
|
||||
|
||||
childrenField.setAccessible(true);
|
||||
taskQueueField.setAccessible(true);
|
||||
tailTasksField.setAccessible(true);
|
||||
}
|
||||
catch (NoSuchFieldException e)
|
||||
{
|
||||
throw new IllegalStateException(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -240,10 +240,12 @@ public final class SSLFactory
|
|||
* Get a netty {@link SslContext} instance.
|
||||
*/
|
||||
@VisibleForTesting
|
||||
static SslContext getOrCreateSslContext(EncryptionOptions options, boolean buildTruststore,
|
||||
SocketType socketType, boolean useOpenSsl) throws IOException
|
||||
static SslContext getOrCreateSslContext(EncryptionOptions options,
|
||||
boolean buildTruststore,
|
||||
SocketType socketType,
|
||||
boolean useOpenSsl) throws IOException
|
||||
{
|
||||
CacheKey key = new CacheKey(options, socketType);
|
||||
CacheKey key = new CacheKey(options, socketType, useOpenSsl);
|
||||
SslContext sslContext;
|
||||
|
||||
sslContext = cachedSslContexts.get(key);
|
||||
|
|
@ -413,11 +415,13 @@ public final class SSLFactory
|
|||
{
|
||||
private final EncryptionOptions encryptionOptions;
|
||||
private final SocketType socketType;
|
||||
private final boolean useOpenSSL;
|
||||
|
||||
public CacheKey(EncryptionOptions encryptionOptions, SocketType socketType)
|
||||
public CacheKey(EncryptionOptions encryptionOptions, SocketType socketType, boolean useOpenSSL)
|
||||
{
|
||||
this.encryptionOptions = encryptionOptions;
|
||||
this.socketType = socketType;
|
||||
this.useOpenSSL = useOpenSSL;
|
||||
}
|
||||
|
||||
public boolean equals(Object o)
|
||||
|
|
@ -426,6 +430,7 @@ public final class SSLFactory
|
|||
if (o == null || getClass() != o.getClass()) return false;
|
||||
CacheKey cacheKey = (CacheKey) o;
|
||||
return (socketType == cacheKey.socketType &&
|
||||
useOpenSSL == cacheKey.useOpenSSL &&
|
||||
Objects.equals(encryptionOptions, cacheKey.encryptionOptions));
|
||||
}
|
||||
|
||||
|
|
@ -434,6 +439,7 @@ public final class SSLFactory
|
|||
int result = 0;
|
||||
result += 31 * socketType.hashCode();
|
||||
result += 31 * encryptionOptions.hashCode();
|
||||
result += 31 * Boolean.hashCode(useOpenSSL);
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -153,7 +153,7 @@ public class NativeTransportService
|
|||
final boolean enableEpoll = Boolean.parseBoolean(System.getProperty("cassandra.native.epoll.enabled", "true"));
|
||||
|
||||
if (enableEpoll && !Epoll.isAvailable() && NativeLibrary.osType == NativeLibrary.OSType.LINUX)
|
||||
logger.warn("epoll not available {}", Epoll.unavailabilityCause());
|
||||
logger.warn("epoll not available", Epoll.unavailabilityCause());
|
||||
|
||||
return enableEpoll && Epoll.isAvailable();
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue