Skip to content

Commit 8f2fc5b

Browse files
j-bahrmeta-codesync[bot]
authored andcommitted
Plumb NiftyMetrics into transport layer to allow exporting of netty counters
Summary: SPINiftyMetrics provides low level counters for netty transport layer events such as channel count, bytes written etc. The metrics class was new’ed up low in the transport factory layer but not surfaces to service framework in a way that would allow the counters to be published to fb303. This was was a sev followup task T234357428 to make these available as additonal way to diagnose conneciton issues. Reviewed By: RayanRal, adolfojunior Differential Revision: D89578965 fbshipit-source-id: d7b5f99ca1fb133d86953c328cd0adc593d6fe92
1 parent ade9f96 commit 8f2fc5b

11 files changed

Lines changed: 53 additions & 22 deletions

File tree

third-party/thrift/src/thrift/lib/java/benchmarks/src/main/java/com/facebook/thrift/jmh/ReactiveRpcBenchmarks.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import com.facebook.thrift.example.ping.PingServiceReactiveClient;
2626
import com.facebook.thrift.example.ping.PingServiceRpcServerHandler;
2727
import com.facebook.thrift.legacy.server.LegacyServerTransportFactory;
28+
import com.facebook.thrift.util.SPINiftyMetrics;
2829
import com.facebook.thrift.util.resources.RpcResources;
2930
import io.netty.channel.unix.DomainSocketAddress;
3031
import java.net.InetSocketAddress;
@@ -84,7 +85,7 @@ public void setup(Blackhole bh) {
8485
LegacyServerTransportFactory factory =
8586
new LegacyServerTransportFactory(
8687
new ThriftServerConfig().setSslEnabled(false).setUdsPath("/tmp/jmh.socket"));
87-
factory.createServerTransport(socketAddress, serverHandler).block();
88+
factory.createServerTransport(socketAddress, serverHandler, new SPINiftyMetrics()).block();
8889
System.out.println("Ping Service started...");
8990

9091
System.out.println("Connecting Ping Service Client...");

third-party/thrift/src/thrift/lib/java/benchmarks/src/main/java/com/facebook/thrift/runner/MultiUdsReactiveServer.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import com.facebook.thrift.example.ping.PingServiceRpcServerHandler;
2424
import com.facebook.thrift.legacy.server.LegacyServerTransport;
2525
import com.facebook.thrift.legacy.server.LegacyServerTransportFactory;
26+
import com.facebook.thrift.util.SPINiftyMetrics;
2627
import io.netty.channel.unix.DomainSocketAddress;
2728
import io.netty.util.ResourceLeakDetector;
2829
import java.util.ArrayList;
@@ -62,7 +63,9 @@ public static void main(String... args) {
6263
.setUdsPath("/tmp/uds_benchmark.socket" + ix));
6364
DomainSocketAddress socketAddress = new DomainSocketAddress("/tmp/uds_benchmark.socket" + ix);
6465
LegacyServerTransport transport =
65-
transportFactory.createServerTransport(socketAddress, serverHandler).block();
66+
transportFactory
67+
.createServerTransport(socketAddress, serverHandler, new SPINiftyMetrics())
68+
.block();
6669
System.out.println("creating server listen on uds: " + socketAddress);
6770
transports.add(transport);
6871
}

third-party/thrift/src/thrift/lib/java/benchmarks/src/main/java/com/facebook/thrift/runner/UdsReactiveServer.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import com.facebook.thrift.example.ping.PingServiceRpcServerHandler;
2424
import com.facebook.thrift.legacy.server.LegacyServerTransport;
2525
import com.facebook.thrift.legacy.server.LegacyServerTransportFactory;
26+
import com.facebook.thrift.util.SPINiftyMetrics;
2627
import io.netty.channel.unix.DomainSocketAddress;
2728
import io.netty.util.ResourceLeakDetector;
2829
import java.util.Collections;
@@ -48,7 +49,9 @@ public static void main(String... args) {
4849
new ThriftServerConfig().setSslEnabled(false).setUdsPath("/tmp/uds_benchmark.socket"));
4950
DomainSocketAddress socketAddress = new DomainSocketAddress("/tmp/uds_benchmark.socket");
5051
LegacyServerTransport transport =
51-
transportFactory.createServerTransport(socketAddress, serverHandler).block();
52+
transportFactory
53+
.createServerTransport(socketAddress, serverHandler, new SPINiftyMetrics())
54+
.block();
5255

5356
LockSupport.park();
5457
}

third-party/thrift/src/thrift/lib/java/runtime/src/main/java/com/facebook/swift/service/stats/ServerStats.java

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,9 @@ public class ServerStats {
5656
private static final String ACTIVE_REQUESTS_KEY = "thrift.active_requests.avg";
5757
private static final String QUEUED_REQUESTS_KEY = "thrift.queued_requests.avg.60";
5858
// Connection related Counter Keys
59+
private static final String CHANNEL_COUNT_KEY = "thrift.channel.count";
60+
private static final String BYTES_READ_KEY = "thrift.bytes_read.count";
61+
private static final String BYTES_WRITTEN_KEY = "thrift.bytes_written.count";
5962
private static final String ACCEPTED_CONNS_KEY = "thrift.accepted_connections.count";
6063
private static final String DROPPED_CONNS_KEY = "thrift.dropped_conns.count";
6164
private static final String REJECTED_CONNS_KEY = "thrift.rejected_conns.count";
@@ -146,6 +149,11 @@ public Map<String, Long> getCounters() {
146149
resultCounters.put(QUEUED_REQUESTS_KEY, (long) threadPoolExecutor.getQueue().size());
147150
}
148151
if (niftyMetrics != null) {
152+
resultCounters.put(CHANNEL_COUNT_KEY, (long) niftyMetrics.getChannelCount());
153+
154+
resultCounters.put(BYTES_READ_KEY, niftyMetrics.getBytesRead());
155+
resultCounters.put(BYTES_WRITTEN_KEY, niftyMetrics.getBytesWritten());
156+
149157
resultCounters.put(ACCEPTED_CONNS_KEY, niftyMetrics.getAcceptedConnections());
150158
resultCounters.put(
151159
ACCEPTED_CONNS_KEY + ONE_MINUTE, niftyMetrics.getAcceptedConnectionsOneMin());

third-party/thrift/src/thrift/lib/java/runtime/src/main/java/com/facebook/thrift/legacy/server/LegacyServerTransport.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -57,15 +57,17 @@ public class LegacyServerTransport implements ServerTransport {
5757
}
5858

5959
static Mono<LegacyServerTransport> createNewInstance(
60-
SocketAddress bindAddress, RpcServerHandler rpcServerHandler, ThriftServerConfig config) {
60+
SocketAddress bindAddress,
61+
RpcServerHandler rpcServerHandler,
62+
ThriftServerConfig config,
63+
SPINiftyMetrics serverMetrics) {
6164

6265
return Mono.defer(
6366
() -> {
6467
requireNonNull(rpcServerHandler, "methodInvoker is null");
6568
requireNonNull(config, "config is null");
6669

6770
MonoProcessor<Void> onClose = MonoProcessor.create();
68-
SPINiftyMetrics metrics = new SPINiftyMetrics();
6971

7072
Optional<Supplier<SslContext>> sslContext = Optional.empty();
7173
if (config.isSslEnabled() && bindAddress instanceof InetSocketAddress) {
@@ -80,7 +82,7 @@ static Mono<LegacyServerTransport> createNewInstance(
8082
sslContext,
8183
config.isAllowPlaintext(),
8284
config.isAssumeClientsSupportOutOfOrderResponses(),
83-
metrics,
85+
serverMetrics,
8486
config.getConnectionLimit());
8587

8688
EventLoopGroup group = RpcResources.getEventLoopGroup();
@@ -116,7 +118,7 @@ static Mono<LegacyServerTransport> createNewInstance(
116118
NettyUtil.toMono(channel.closeFuture()).subscribe(onClose);
117119

118120
return NettyUtil.toMono(bind)
119-
.thenReturn(new LegacyServerTransport(channel, metrics, onClose));
121+
.thenReturn(new LegacyServerTransport(channel, serverMetrics, onClose));
120122
});
121123
}
122124

third-party/thrift/src/thrift/lib/java/runtime/src/main/java/com/facebook/thrift/legacy/server/LegacyServerTransportFactory.java

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
import com.facebook.swift.service.ThriftServerConfig;
2020
import com.facebook.thrift.server.RpcServerHandler;
2121
import com.facebook.thrift.server.ServerTransportFactory;
22+
import com.facebook.thrift.util.SPINiftyMetrics;
2223
import java.net.InetSocketAddress;
2324
import java.net.SocketAddress;
2425
import reactor.core.publisher.Mono;
@@ -32,13 +33,16 @@ public LegacyServerTransportFactory(ThriftServerConfig config) {
3233

3334
@Override
3435
public Mono<? extends LegacyServerTransport> createServerTransport(
35-
SocketAddress bindAddress, RpcServerHandler rpcServerHandler) {
36-
return LegacyServerTransport.createNewInstance(bindAddress, rpcServerHandler, config);
36+
SocketAddress bindAddress, RpcServerHandler rpcServerHandler, SPINiftyMetrics serverMetrics) {
37+
return LegacyServerTransport.createNewInstance(
38+
bindAddress, rpcServerHandler, config, serverMetrics);
3739
}
3840

3941
public Mono<? extends LegacyServerTransport> createServerTransport(
4042
RpcServerHandler rpcServerHandler) {
4143
return createServerTransport(
42-
new InetSocketAddress("localhost", config.getPort()), rpcServerHandler);
44+
new InetSocketAddress("localhost", config.getPort()),
45+
rpcServerHandler,
46+
new SPINiftyMetrics());
4347
}
4448
}

third-party/thrift/src/thrift/lib/java/runtime/src/main/java/com/facebook/thrift/rsocket/server/RSocketServerTransport.java

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -44,17 +44,18 @@ public class RSocketServerTransport implements ServerTransport {
4444
}
4545

4646
static Mono<RSocketServerTransport> createInstance(
47-
SocketAddress socketAddress, RpcServerHandler rpcServerHandler, ThriftServerConfig config) {
47+
SocketAddress socketAddress,
48+
RpcServerHandler rpcServerHandler,
49+
ThriftServerConfig config,
50+
SPINiftyMetrics serverMetrics) {
4851
try {
4952
requireNonNull(rpcServerHandler, "methodInvoker is null");
5053
requireNonNull(config, "config is null");
5154

52-
SPINiftyMetrics metrics = new SPINiftyMetrics();
53-
5455
return RSocketServer.create(new ThriftSocketAcceptor(rpcServerHandler))
5556
.fragment(MAX_FRAME_SIZE)
5657
.payloadDecoder(PayloadDecoder.ZERO_COPY)
57-
.bind(new ReactorServerTransport(socketAddress, config, metrics))
58+
.bind(new ReactorServerTransport(socketAddress, config, serverMetrics))
5859
.map(RSocketServerTransport::new);
5960
} catch (Exception e) {
6061
return Mono.error(e);

third-party/thrift/src/thrift/lib/java/runtime/src/main/java/com/facebook/thrift/rsocket/server/RSocketServerTransportFactory.java

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import com.facebook.thrift.server.RpcServerHandler;
2121
import com.facebook.thrift.server.ServerTransportFactory;
2222
import com.facebook.thrift.util.RpcServerUtils;
23+
import com.facebook.thrift.util.SPINiftyMetrics;
2324
import java.net.InetSocketAddress;
2425
import java.net.SocketAddress;
2526
import reactor.core.publisher.Mono;
@@ -34,14 +35,17 @@ public RSocketServerTransportFactory(ThriftServerConfig config) {
3435

3536
@Override
3637
public Mono<? extends RSocketServerTransport> createServerTransport(
37-
SocketAddress bindAddress, RpcServerHandler rpcServerHandler) {
38-
return RSocketServerTransport.createInstance(parsePort(bindAddress), rpcServerHandler, config);
38+
SocketAddress bindAddress, RpcServerHandler rpcServerHandler, SPINiftyMetrics serverMetrics) {
39+
return RSocketServerTransport.createInstance(
40+
parsePort(bindAddress), rpcServerHandler, config, serverMetrics);
3941
}
4042

4143
public Mono<? extends RSocketServerTransport> createServerTransport(
4244
RpcServerHandler rpcServerHandler) {
4345
return createServerTransport(
44-
new InetSocketAddress("localhost", parsePort(config.getPort())), rpcServerHandler);
46+
new InetSocketAddress("localhost", parsePort(config.getPort())),
47+
rpcServerHandler,
48+
new SPINiftyMetrics());
4549
}
4650

4751
private int parsePort(int port) {

third-party/thrift/src/thrift/lib/java/runtime/src/main/java/com/facebook/thrift/server/ServerTransportFactory.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,10 +16,11 @@
1616

1717
package com.facebook.thrift.server;
1818

19+
import com.facebook.thrift.util.SPINiftyMetrics;
1920
import java.net.SocketAddress;
2021
import reactor.core.publisher.Mono;
2122

2223
public interface ServerTransportFactory<T extends ServerTransport> {
2324
Mono<? extends T> createServerTransport(
24-
SocketAddress bindAddress, RpcServerHandler rpcServerHandler);
25+
SocketAddress bindAddress, RpcServerHandler rpcServerHandler, SPINiftyMetrics niftyMetrics);
2526
}

third-party/thrift/src/thrift/lib/java/runtime/src/main/java/com/facebook/thrift/util/RpcServerUtils.java

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -216,15 +216,16 @@ public static Mono<? extends ServerTransport> createServerTransport(
216216

217217
final int port = config.getPort() == 0 ? RpcServerUtils.findFreePort() : config.getPort();
218218

219+
SPINiftyMetrics serverMetrics = new SPINiftyMetrics();
219220
if (config.isEnableUDS()) {
220221
return transportFactory.createServerTransport(
221-
new DomainSocketAddress(config.getUdsPath()), rpcServerHandler);
222+
new DomainSocketAddress(config.getUdsPath()), rpcServerHandler, serverMetrics);
222223
} else if (config.isBindAddressEnabled()) {
223224
return transportFactory.createServerTransport(
224-
new InetSocketAddress(config.getBindAddress(), port), rpcServerHandler);
225+
new InetSocketAddress(config.getBindAddress(), port), rpcServerHandler, serverMetrics);
225226
} else {
226227
return transportFactory.createServerTransport(
227-
new InetSocketAddress("localhost", port), rpcServerHandler);
228+
new InetSocketAddress("localhost", port), rpcServerHandler, serverMetrics);
228229
}
229230
}
230231

0 commit comments

Comments
 (0)