Skip to content

Commit ede36cd

Browse files
NIFI-16152 Switched from Jersey to Jetty for gzip compression (#11489)
- Added drain and close for replicated request response body to handle gzip encoded responses
1 parent 1f6d195 commit ede36cd

19 files changed

Lines changed: 521 additions & 58 deletions

File tree

nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/http/replication/StandardAsyncClusterResponse.java

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717

1818
package org.apache.nifi.cluster.coordination.http.replication;
1919

20+
import jakarta.ws.rs.core.Response;
2021
import org.apache.nifi.cluster.coordination.http.HttpResponseMapper;
2122
import org.apache.nifi.cluster.manager.NodeResponse;
2223
import org.apache.nifi.cluster.protocol.NodeIdentifier;
@@ -216,6 +217,7 @@ public synchronized NodeResponse getMergedResponse(final boolean triggerCallback
216217

217218
final long start = System.nanoTime();
218219
mergedResponse = responseMapper.mapResponses(uri, method, nodeResponses, merge);
220+
closeUnusedResponses(nodeResponses);
219221
final long nanos = System.nanoTime() - start;
220222
addTiming("Map/Merge Responses", "All Nodes", nanos);
221223

@@ -329,6 +331,23 @@ public String toString() {
329331
+ ", responses=" + getCompletedNodeIdentifiers().size() + "/" + responseMap.size() + "]";
330332
}
331333

334+
private void closeUnusedResponses(final Set<NodeResponse> nodeResponses) {
335+
final Response clientMergedResponse = mergedResponse.getClientResponse();
336+
if (mergedResponse.getUpdatedEntity() != null) {
337+
// Close merged response since Updated Entity used in place of stream
338+
clientMergedResponse.close();
339+
}
340+
341+
for (final NodeResponse nodeResponse : nodeResponses) {
342+
final Response clientNodeResponse = nodeResponse.getClientResponse();
343+
if (clientMergedResponse == clientNodeResponse) {
344+
continue;
345+
}
346+
347+
nodeResponse.close();
348+
}
349+
}
350+
332351
private static class ResponseHolder {
333352
private final long nanoStart;
334353
private long requestNanos;

nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/http/replication/ThreadPoolRequestReplicator.java

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -523,8 +523,13 @@ public void onCompletion(final NodeResponse nodeResponse) {
523523
// If all nodes responded with 202-Accepted, then we can replicate the original request
524524
// to all nodes and we are finished.
525525
if (dissentingCount == 0) {
526-
logger.debug("Received verification from all {} nodes that mutable request {} {} can be made", numNodes, method, uri.getPath());
527-
replicate(nodeIds, method, uri, entity, headers, false, clusterResponse, true, merge, monitor);
526+
try {
527+
logger.debug("Received verification from all {} nodes that mutable request {} {} can be made", numNodes, method, uri.getPath());
528+
replicate(nodeIds, method, uri, entity, headers, false, clusterResponse, true, merge, monitor);
529+
} finally {
530+
// Close HTTP Responses after replication completed
531+
nodeResponses.forEach(NodeResponse::close);
532+
}
528533
return;
529534
}
530535

@@ -582,6 +587,9 @@ public void onCompletion(final NodeResponse nodeResponse) {
582587
clusterResponse.setFailure(failure, response.getNodeId());
583588
}
584589
}
590+
591+
// Close verification responses for cancelled transactions
592+
nodeResponses.forEach(NodeResponse::close);
585593
} finally {
586594
if (monitor != null) {
587595
synchronized (monitor) {

nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/http/replication/client/StandardHttpReplicationClient.java

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package org.apache.nifi.cluster.coordination.http.replication.client;
1818

1919
import com.fasterxml.jackson.annotation.JsonInclude;
20+
import com.fasterxml.jackson.core.JsonParser;
2021
import com.fasterxml.jackson.databind.ObjectMapper;
2122
import com.fasterxml.jackson.module.jakarta.xmlbind.JakartaXmlBindAnnotationIntrospector;
2223
import jakarta.ws.rs.core.MultivaluedHashMap;
@@ -41,6 +42,7 @@
4142
import java.io.ByteArrayOutputStream;
4243
import java.io.IOException;
4344
import java.io.InputStream;
45+
import java.io.OutputStream;
4446
import java.io.UncheckedIOException;
4547
import java.net.URI;
4648
import java.util.LinkedHashMap;
@@ -113,6 +115,8 @@ public StandardHttpReplicationClient(final WebClientService webClientService, fi
113115

114116
objectMapper.setDefaultPropertyInclusion(JsonInclude.Value.construct(JsonInclude.Include.NON_NULL, JsonInclude.Include.ALWAYS));
115117
objectMapper.setAnnotationIntrospector(new JakartaXmlBindAnnotationIntrospector(objectMapper.getTypeFactory()));
118+
// Disable closing source streams to allow draining and subsequent closing
119+
objectMapper.disable(JsonParser.Feature.AUTO_CLOSE_SOURCE);
116120

117121
jsonSerializer = new JsonEntitySerializer(objectMapper);
118122
xmlSerializer = new XmlEntitySerializer();
@@ -208,11 +212,25 @@ private Response replicate(final HttpRequestBodySpec httpRequestBodySpec, final
208212

209213
final InputStream responseBody = getResponseBody(responseEntity.body(), headers);
210214
final Runnable closeCallback = () -> {
215+
try {
216+
// Drain raw response stream before closing the Response Entity
217+
responseEntity.body().transferTo(OutputStream.nullOutputStream());
218+
} catch (final IOException e) {
219+
logger.debug("Drain failed for Replicated {} {} HTTP {}", method, location, statusCode, e);
220+
}
221+
211222
try {
212223
responseEntity.close();
213224
} catch (final IOException e) {
214225
logger.warn("Close failed for Replicated {} {} HTTP {}", method, location, statusCode, e);
215226
}
227+
228+
try {
229+
// Release resources for gzip wrapped streams
230+
responseBody.close();
231+
} catch (final IOException e) {
232+
logger.warn("Close failed for Replicated Response Body {} {} HTTP {}", method, location, statusCode, e);
233+
}
216234
};
217235

218236
final long elapsed = System.currentTimeMillis() - started;

nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/http/replication/io/ReplicatedResponse.java

Lines changed: 39 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -17,8 +17,6 @@
1717

1818
package org.apache.nifi.cluster.coordination.http.replication.io;
1919

20-
import com.fasterxml.jackson.core.JsonFactory;
21-
import com.fasterxml.jackson.core.JsonParser;
2220
import com.fasterxml.jackson.databind.ObjectMapper;
2321
import jakarta.ws.rs.core.EntityTag;
2422
import jakarta.ws.rs.core.GenericType;
@@ -51,28 +49,28 @@ public class ReplicatedResponse extends Response {
5149
private static final int MAXIMUM_BUFFER_SIZE = 1048576;
5250
private static final int CONTENT_LENGTH_UNKNOWN = -1;
5351

54-
private final ObjectMapper codec;
52+
private final ObjectMapper objectMapper;
5553
private final InputStream responseBody;
5654
private final MultivaluedMap<String, String> responseHeaders;
5755
private final URI location;
5856
private final int statusCode;
5957
private final Runnable closeCallback;
6058
private final int contentLength;
6159

62-
private final JsonFactory jsonFactory = new JsonFactory();
63-
6460
private final byte[] bufferedResponseBody;
6561

62+
private Object bufferedEntity;
63+
6664
public ReplicatedResponse(
67-
final ObjectMapper codec,
65+
final ObjectMapper objectMapper,
6866
final InputStream responseBody,
6967
final MultivaluedMap<String, String> responseHeaders,
7068
final URI location,
7169
final int statusCode,
7270
final int contentLength,
7371
final Runnable closeCallback
7472
) {
75-
this.codec = codec;
73+
this.objectMapper = objectMapper;
7674
this.responseBody = responseBody;
7775
this.responseHeaders = responseHeaders;
7876
this.location = location;
@@ -101,42 +99,30 @@ public StatusType getStatusInfo() {
10199

102100
@Override
103101
public Object getEntity() {
104-
final InputStream responseBodyStream = getResponseBodyStream();
105-
106-
try {
107-
final JsonParser parser = jsonFactory.createParser(responseBodyStream);
108-
parser.setCodec(codec);
109-
return parser.readValueAs(Object.class);
110-
} catch (final Exception e) {
111-
throw new RuntimeException("Failed to parse response", e);
102+
if (bufferedEntity == null) {
103+
// Read response entity to buffered entity to support multiple invocations
104+
bufferedEntity = readResponseEntity(Object.class);
112105
}
106+
107+
return bufferedEntity;
113108
}
114109

115110
@Override
116111
@SuppressWarnings("unchecked")
117-
public <T> T readEntity(Class<T> entityType) {
118-
final InputStream responseBodyStream = getResponseBodyStream();
112+
public <T> T readEntity(final Class<T> entityType) {
113+
final T entity;
119114

115+
// Return raw response body stream when requested without buffering
120116
if (InputStream.class.equals(entityType)) {
121-
return (T) responseBodyStream;
122-
}
123-
124-
if (String.class.equals(entityType)) {
125-
try {
126-
final byte[] responseBytes = responseBodyStream.readAllBytes();
127-
return (T) new String(responseBytes, StandardCharsets.UTF_8);
128-
} catch (final IOException e) {
129-
throw new UncheckedIOException("Read Replicated Response Body to String failed for %s".formatted(location), e);
130-
}
117+
return (T) getResponseBodyStream();
131118
}
132119

133-
try {
134-
final JsonParser parser = jsonFactory.createParser(responseBodyStream);
135-
parser.setCodec(codec);
136-
return parser.readValueAs(entityType);
137-
} catch (final Exception e) {
138-
throw new RuntimeException("Failed to parse response as entity of type " + entityType, e);
120+
if (bufferedEntity == null) {
121+
// Read response entity to buffered entity to support multiple invocations
122+
bufferedEntity = readResponseEntity(entityType);
139123
}
124+
entity = (T) bufferedEntity;
125+
return entity;
140126
}
141127

142128
@Override
@@ -287,6 +273,26 @@ private InputStream getResponseBodyStream() {
287273
return responseBodyStream;
288274
}
289275

276+
@SuppressWarnings("unchecked")
277+
private <T> T readResponseEntity(final Class<T> entityType) {
278+
final InputStream responseBodyStream = getResponseBodyStream();
279+
280+
if (String.class.equals(entityType)) {
281+
try {
282+
final byte[] responseBytes = responseBodyStream.readAllBytes();
283+
return (T) new String(responseBytes, StandardCharsets.UTF_8);
284+
} catch (final IOException e) {
285+
throw new UncheckedIOException("Read Replicated Response Body to String failed for %s".formatted(location), e);
286+
}
287+
}
288+
289+
try {
290+
return objectMapper.readValue(responseBodyStream, entityType);
291+
} catch (final Exception e) {
292+
throw new RuntimeException("Failed to parse response as Entity [%s] for %s".formatted(entityType, location), e);
293+
}
294+
}
295+
290296
private static byte[] readResponseBody(final InputStream inputStream, final URI location, final int statusCode) {
291297
try {
292298
return inputStream.readAllBytes();

nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/manager/NodeResponse.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import org.slf4j.Logger;
3030
import org.slf4j.LoggerFactory;
3131

32+
import java.io.Closeable;
3233
import java.io.InputStream;
3334
import java.net.URI;
3435
import java.util.List;
@@ -47,7 +48,7 @@
4748
* This class overrides hashCode and equals and considers two instances to be equal if they have the equal NodeIdentifiers.
4849
*
4950
*/
50-
public class NodeResponse {
51+
public class NodeResponse implements Closeable {
5152

5253
private static final Logger logger = LoggerFactory.getLogger(NodeResponse.class);
5354
private final String httpMethod;
@@ -269,4 +270,11 @@ public String toString() {
269270
.append(",Duration=").append(TimeUnit.MILLISECONDS.convert(requestDurationNanos, TimeUnit.NANOSECONDS)).append(" ms]");
270271
return sb.toString();
271272
}
273+
274+
@Override
275+
public void close() {
276+
if (response != null) {
277+
response.close();
278+
}
279+
}
272280
}

nifi-framework-bundle/nifi-framework/nifi-web/nifi-jetty/pom.xml

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,14 @@
8585
<artifactId>nifi-web-servlet-shared</artifactId>
8686
<version>2.11.0-SNAPSHOT</version>
8787
</dependency>
88+
<dependency>
89+
<groupId>org.eclipse.jetty.compression</groupId>
90+
<artifactId>jetty-compression-server</artifactId>
91+
</dependency>
92+
<dependency>
93+
<groupId>org.eclipse.jetty.compression</groupId>
94+
<artifactId>jetty-compression-gzip</artifactId>
95+
</dependency>
8896
<dependency>
8997
<groupId>org.eclipse.jetty</groupId>
9098
<artifactId>jetty-deploy</artifactId>

nifi-framework-bundle/nifi-framework/nifi-web/nifi-jetty/src/main/java/org/apache/nifi/web/server/JettyServer.java

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@
6969
import org.apache.nifi.web.server.filter.RequestFilterProvider;
7070
import org.apache.nifi.web.server.filter.RestApiRequestFilterProvider;
7171
import org.apache.nifi.web.server.filter.StandardRequestFilterProvider;
72+
import org.eclipse.jetty.compression.server.CompressionHandler;
7273
import org.eclipse.jetty.deploy.StandardDeployer;
7374
import org.eclipse.jetty.ee.webapp.WebAppClassLoader;
7475
import org.eclipse.jetty.ee11.servlet.ErrorPageErrorHandler;
@@ -222,7 +223,13 @@ public void init() {
222223
deployer = new StandardDeployer(contextHandlerCollection);
223224
server.addBean(deployer);
224225

225-
serverHandlerCollection.addHandler(contextHandlerCollection);
226+
// Deploy the web applications beneath the CompressionHandler
227+
final CompressionHandler compressionHandler = serverHandlerCollection.getDescendant(CompressionHandler.class);
228+
if (compressionHandler == null) {
229+
throw new IllegalStateException("Compression Handler not configured: Server Provider configuration failed");
230+
} else {
231+
compressionHandler.setHandler(contextHandlerCollection);
232+
}
226233
} else {
227234
throw new IllegalStateException("Server Handler not Handler.Collection: Server Provider configuration failed");
228235
}

nifi-framework-bundle/nifi-framework/nifi-web/nifi-jetty/src/main/java/org/apache/nifi/web/server/StandardServerProvider.java

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,8 +21,12 @@
2121
import org.apache.nifi.web.server.connector.FrameworkServerConnectorFactory;
2222
import org.apache.nifi.web.server.handler.ContextPathRedirectPatternRule;
2323
import org.apache.nifi.web.server.handler.HeaderWriterHandler;
24+
import org.apache.nifi.web.server.handler.UnsupportedContentEncodingHandler;
2425
import org.apache.nifi.web.server.log.RequestLogProvider;
2526
import org.apache.nifi.web.server.log.StandardRequestLogProvider;
27+
import org.eclipse.jetty.compression.gzip.GzipCompression;
28+
import org.eclipse.jetty.compression.server.CompressionConfig;
29+
import org.eclipse.jetty.compression.server.CompressionHandler;
2630
import org.eclipse.jetty.rewrite.handler.RedirectPatternRule;
2731
import org.eclipse.jetty.rewrite.handler.RewriteHandler;
2832
import org.eclipse.jetty.server.Handler;
@@ -47,6 +51,8 @@
4751
* Standard implementation of Server Provider with default Handlers
4852
*/
4953
class StandardServerProvider implements ServerProvider {
54+
private static final String ROOT_PATH = "/";
55+
5056
private static final String ALL_PATHS_PATTERN = "/*";
5157

5258
private static final String FRONTEND_CONTEXT_PATH = "/nifi";
@@ -139,6 +145,27 @@ private Handler getStandardHandler() {
139145
// Set Handler for standard response headers
140146
standardHandler.addHandler(new HeaderWriterHandler());
141147

148+
// Reject requests that declare a Content-Encoding because request bodies are not decompressed
149+
standardHandler.addHandler(new UnsupportedContentEncodingHandler());
150+
151+
// Set Handler for response compression
152+
standardHandler.addHandler(getCompressionHandler());
153+
142154
return standardHandler;
143155
}
156+
157+
private CompressionHandler getCompressionHandler() {
158+
final GzipCompression gzipCompression = new GzipCompression();
159+
final CompressionHandler compressionHandler = new CompressionHandler();
160+
compressionHandler.putCompression(gzipCompression);
161+
162+
final CompressionConfig compressionConfig = CompressionConfig.builder()
163+
.defaults()
164+
// Disable decompression of requests
165+
.decompressExcludeEncoding(gzipCompression.getEncodingName())
166+
.build();
167+
compressionHandler.putConfiguration(ROOT_PATH, compressionConfig);
168+
169+
return compressionHandler;
170+
}
144171
}

0 commit comments

Comments
 (0)