io.grpc
grpc-inprocess
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/CallOptions.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/CallOptions.java
index b6e052d223..dd16de0167 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/CallOptions.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/CallOptions.java
@@ -17,7 +17,11 @@
package org.apache.arrow.flight;
import io.grpc.stub.AbstractStub;
+import java.util.Arrays;
import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import org.apache.arrow.vector.compression.CompressionCodec;
+import org.apache.arrow.vector.compression.CompressionUtil;
/** Common call options. */
public class CallOptions {
@@ -25,6 +29,25 @@ public static CallOption timeout(long duration, TimeUnit unit) {
return new Timeout(duration, unit);
}
+ /**
+ * Advertise support for IPC body compression codecs in preference order.
+ *
+ * Servers that do not support negotiation ignore this option and return uncompressed IPC.
+ */
+ public static CallOption acceptIpcCompression(CompressionUtil.CodecType... codecs) {
+ if (codecs.length == 0) {
+ throw new IllegalArgumentException("At least one IPC compression codec is required");
+ }
+ for (CompressionUtil.CodecType codec : codecs) {
+ CompressionCodec.Factory.INSTANCE.createCodec(codec);
+ }
+ final FlightCallHeaders headers = new FlightCallHeaders();
+ headers.insert(
+ FlightConstants.IPC_ACCEPT_COMPRESSION_HEADER,
+ Arrays.stream(codecs).map(IpcCompression::codecName).collect(Collectors.joining(",")));
+ return new HeaderCallOption(headers);
+ }
+
static > T wrapStub(T stub, CallOption[] options) {
for (CallOption option : options) {
if (option instanceof GrpcCallOption) {
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/DictionaryUtils.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/DictionaryUtils.java
index cecc1b876e..1ca9b424b5 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/DictionaryUtils.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/DictionaryUtils.java
@@ -28,6 +28,7 @@
import org.apache.arrow.vector.FieldVector;
import org.apache.arrow.vector.VectorSchemaRoot;
import org.apache.arrow.vector.VectorUnloader;
+import org.apache.arrow.vector.compression.CompressionCodec;
import org.apache.arrow.vector.dictionary.Dictionary;
import org.apache.arrow.vector.dictionary.DictionaryProvider;
import org.apache.arrow.vector.ipc.message.ArrowDictionaryBatch;
@@ -57,6 +58,18 @@ static Schema generateSchemaMessages(
final IpcOption option,
final Consumer messageCallback)
throws Exception {
+ return generateSchemaMessages(
+ originalSchema, descriptor, provider, option, null, messageCallback);
+ }
+
+ static Schema generateSchemaMessages(
+ final Schema originalSchema,
+ final FlightDescriptor descriptor,
+ final DictionaryProvider provider,
+ final IpcOption option,
+ final CompressionCodec compressionCodec,
+ final Consumer messageCallback)
+ throws Exception {
final Set dictionaryIds = new HashSet<>();
final Schema schema = generateSchema(originalSchema, provider, dictionaryIds);
MetadataV4UnionChecker.checkForUnion(schema.getFields().iterator(), option.metadataVersion);
@@ -77,7 +90,9 @@ static Schema generateSchemaMessages(
Collections.singletonList(vector.getField()),
Collections.singletonList(vector),
count);
- final VectorUnloader unloader = new VectorUnloader(dictRoot);
+ final VectorUnloader unloader =
+ new VectorUnloader(
+ dictRoot, /* includeNullCount */ true, compressionCodec, /* alignBuffers */ true);
try (final ArrowDictionaryBatch dictionaryBatch =
new ArrowDictionaryBatch(id, unloader.getRecordBatch());
final ArrowMessage message = new ArrowMessage(dictionaryBatch, option)) {
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightBindingService.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightBindingService.java
index b68f3aa86c..982e2cc75f 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightBindingService.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightBindingService.java
@@ -26,6 +26,7 @@
import io.grpc.protobuf.ProtoUtils;
import io.grpc.stub.ServerCalls;
import io.grpc.stub.StreamObserver;
+import java.util.List;
import java.util.Set;
import java.util.concurrent.ExecutorService;
import org.apache.arrow.flight.auth.ServerAuthHandler;
@@ -33,6 +34,7 @@
import org.apache.arrow.flight.impl.Flight.PutResult;
import org.apache.arrow.flight.impl.FlightServiceGrpc;
import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.vector.compression.CompressionUtil;
/** Extends the basic flight service to override some methods for more efficient implementations. */
class FlightBindingService implements BindableService {
@@ -53,8 +55,18 @@ public FlightBindingService(
FlightProducer producer,
ServerAuthHandler authHandler,
ExecutorService executor) {
+ this(allocator, producer, authHandler, executor, java.util.Collections.emptyList());
+ }
+
+ public FlightBindingService(
+ BufferAllocator allocator,
+ FlightProducer producer,
+ ServerAuthHandler authHandler,
+ ExecutorService executor,
+ List ipcCompressionCodecs) {
this.allocator = allocator;
- this.delegate = new FlightService(allocator, producer, authHandler, executor);
+ this.delegate =
+ new FlightService(allocator, producer, authHandler, executor, ipcCompressionCodecs);
}
public static MethodDescriptor getDoGetDescriptor(
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightConstants.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightConstants.java
index 6b89c794d6..997a11d950 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightConstants.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightConstants.java
@@ -21,6 +21,9 @@ public interface FlightConstants {
String SERVICE = "arrow.flight.protocol.FlightService";
+ /** Request header listing supported IPC body compression codecs in preference order. */
+ String IPC_ACCEPT_COMPRESSION_HEADER = "arrow-ipc-accept-compression";
+
FlightServerMiddleware.Key HEADER_KEY =
FlightServerMiddleware.Key.of("org.apache.arrow.flight.ServerHeaderMiddleware");
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java
index ac761457f5..ea091a89c4 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightServer.java
@@ -34,6 +34,8 @@
import java.net.URI;
import java.net.URISyntaxException;
import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -55,6 +57,8 @@
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.util.Preconditions;
import org.apache.arrow.util.VisibleForTesting;
+import org.apache.arrow.vector.compression.CompressionCodec;
+import org.apache.arrow.vector.compression.CompressionUtil;
/**
* Generic server of flight data that is customized via construction with delegate classes for the
@@ -197,6 +201,7 @@ public static final class Builder {
private final List> interceptors;
// Keep track of inserted interceptors
private final Set interceptorKeys;
+ private List ipcCompressionCodecs = Collections.emptyList();
Builder() {
builderOptions = new HashMap<>();
@@ -321,7 +326,7 @@ public FlightServer build() {
}
final FlightBindingService flightService =
- new FlightBindingService(allocator, producer, authHandler, exec);
+ new FlightBindingService(allocator, producer, authHandler, exec, ipcCompressionCodecs);
builder
.executor(exec)
.maxInboundMessageSize(maxInboundMessageSize)
@@ -392,6 +397,22 @@ public Builder backpressureThreshold(int backpressureThreshold) {
return this;
}
+ /**
+ * Enable negotiated IPC body compression for server response streams.
+ *
+ * The client preference order wins. Clients that do not advertise support continue to
+ * receive uncompressed IPC.
+ */
+ public Builder ipcCompression(CompressionUtil.CodecType... codecs) {
+ Preconditions.checkArgument(codecs.length > 0, "At least one codec is required");
+ for (CompressionUtil.CodecType codec : codecs) {
+ IpcCompression.codecName(codec);
+ CompressionCodec.Factory.INSTANCE.createCodec(codec);
+ }
+ ipcCompressionCodecs = Collections.unmodifiableList(Arrays.asList(codecs.clone()));
+ return this;
+ }
+
/**
* A small utility function to ensure that InputStream attributes. are closed if they are not
* null
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightService.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightService.java
index 9f130463c0..9382f3b634 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightService.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightService.java
@@ -20,6 +20,7 @@
import io.grpc.stub.ServerCallStreamObserver;
import io.grpc.stub.StreamObserver;
import java.util.Collections;
+import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
@@ -38,6 +39,8 @@
import org.apache.arrow.flight.impl.FlightServiceGrpc.FlightServiceImplBase;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.util.AutoCloseables;
+import org.apache.arrow.vector.compression.CompressionCodec;
+import org.apache.arrow.vector.compression.CompressionUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -51,16 +54,27 @@ class FlightService extends FlightServiceImplBase {
private final FlightProducer producer;
private final ServerAuthHandler authHandler;
private final ExecutorService executors;
+ private final List ipcCompressionCodecs;
FlightService(
BufferAllocator allocator,
FlightProducer producer,
ServerAuthHandler authHandler,
ExecutorService executors) {
+ this(allocator, producer, authHandler, executors, Collections.emptyList());
+ }
+
+ FlightService(
+ BufferAllocator allocator,
+ FlightProducer producer,
+ ServerAuthHandler authHandler,
+ ExecutorService executors,
+ List ipcCompressionCodecs) {
this.allocator = allocator;
this.producer = producer;
this.authHandler = authHandler;
this.executors = new ContextPropagatingExecutorService(executors);
+ this.ipcCompressionCodecs = ipcCompressionCodecs;
}
private CallContext makeContext(ServerCallStreamObserver> responseObserver) {
@@ -107,10 +121,14 @@ public void doGetCustom(
final ServerCallStreamObserver responseObserver =
(ServerCallStreamObserver) responseObserverSimple;
+ final CallContext context = makeContext(responseObserver);
final GetListener listener =
- new GetListener(responseObserver, this::handleExceptionWithMiddleware);
+ new GetListener(
+ responseObserver,
+ this::handleExceptionWithMiddleware,
+ IpcCompression.negotiate(context, ipcCompressionCodecs));
try {
- producer.getStream(makeContext(responseObserver), new Ticket(ticket), listener);
+ producer.getStream(context, new Ticket(ticket), listener);
} catch (Exception ex) {
listener.error(ex);
}
@@ -155,7 +173,14 @@ private static class GetListener extends OutboundStreamListenerImpl
public GetListener(
ServerCallStreamObserver responseObserver, Consumer errorHandler) {
- super(null, responseObserver);
+ this(responseObserver, errorHandler, null);
+ }
+
+ public GetListener(
+ ServerCallStreamObserver responseObserver,
+ Consumer errorHandler,
+ CompressionCodec compressionCodec) {
+ super(null, responseObserver, compressionCodec);
this.errorHandler = errorHandler;
this.completed = false;
this.serverCallResponseObserver = responseObserver;
@@ -327,7 +352,14 @@ private static class ExchangeListener extends GetListener {
public ExchangeListener(
ServerCallStreamObserver responseObserver, Consumer errorHandler) {
- super(responseObserver, errorHandler);
+ this(responseObserver, errorHandler, null);
+ }
+
+ public ExchangeListener(
+ ServerCallStreamObserver responseObserver,
+ Consumer errorHandler,
+ CompressionCodec compressionCodec) {
+ super(responseObserver, errorHandler, compressionCodec);
this.resource = null;
super.setOnCancelHandler(
() -> {
@@ -387,8 +419,12 @@ public StreamObserver doExchangeCustom(
StreamObserver responseObserverSimple) {
final ServerCallStreamObserver responseObserver =
(ServerCallStreamObserver) responseObserverSimple;
+ final CallContext context = makeContext(responseObserver);
final ExchangeListener listener =
- new ExchangeListener(responseObserver, this::handleExceptionWithMiddleware);
+ new ExchangeListener(
+ responseObserver,
+ this::handleExceptionWithMiddleware,
+ IpcCompression.negotiate(context, ipcCompressionCodecs));
final FlightStream fs =
new FlightStream(
allocator,
@@ -405,7 +441,7 @@ public StreamObserver doExchangeCustom(
executors.submit(
() -> {
try {
- producer.doExchange(makeContext(responseObserver), fs, listener);
+ producer.doExchange(context, fs, listener);
} catch (Exception ex) {
listener.error(ex);
}
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightStream.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightStream.java
index 15cfd6ba85..cdec979e43 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightStream.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/FlightStream.java
@@ -38,6 +38,7 @@
import org.apache.arrow.vector.FieldVector;
import org.apache.arrow.vector.VectorLoader;
import org.apache.arrow.vector.VectorSchemaRoot;
+import org.apache.arrow.vector.compression.CompressionUtil;
import org.apache.arrow.vector.dictionary.Dictionary;
import org.apache.arrow.vector.dictionary.DictionaryProvider;
import org.apache.arrow.vector.ipc.message.ArrowDictionaryBatch;
@@ -86,6 +87,7 @@ public void close() throws Exception {}
private volatile Throwable ex;
private volatile ArrowBuf applicationMetadata = null;
@VisibleForTesting volatile MetadataVersion metadataVersion = null;
+ @VisibleForTesting volatile CompressionUtil.CodecType compressionType = null;
/**
* Constructs a new instance.
@@ -272,6 +274,9 @@ public boolean next() {
// Ensure we have the root
root.get().clear();
try (ArrowRecordBatch arb = msg.asRecordBatch()) {
+ compressionType =
+ CompressionUtil.CodecType.fromCompressionType(
+ arb.getBodyCompression().getCodec());
loader.load(arb);
}
updateMetadata(msg);
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/IpcCompression.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/IpcCompression.java
new file mode 100644
index 0000000000..d38c16eb45
--- /dev/null
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/IpcCompression.java
@@ -0,0 +1,72 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.arrow.flight;
+
+import java.util.List;
+import java.util.Locale;
+import org.apache.arrow.vector.compression.CompressionCodec;
+import org.apache.arrow.vector.compression.CompressionUtil;
+
+/** Utilities for negotiating IPC body compression. */
+final class IpcCompression {
+ private IpcCompression() {}
+
+ static String codecName(CompressionUtil.CodecType codec) {
+ switch (codec) {
+ case LZ4_FRAME:
+ return "lz4_frame";
+ case ZSTD:
+ return "zstd";
+ default:
+ throw new IllegalArgumentException("Unsupported IPC compression codec: " + codec);
+ }
+ }
+
+ static CompressionCodec negotiate(
+ FlightProducer.CallContext context, List supportedCodecs) {
+ if (supportedCodecs.isEmpty()) {
+ return null;
+ }
+ final ServerHeaderMiddleware middleware = context.getMiddleware(FlightConstants.HEADER_KEY);
+ if (middleware == null) {
+ return null;
+ }
+ final String accepted = middleware.headers().get(FlightConstants.IPC_ACCEPT_COMPRESSION_HEADER);
+ if (accepted == null) {
+ return null;
+ }
+ for (String name : accepted.split(",")) {
+ final CompressionUtil.CodecType codec = parseCodec(name);
+ if (codec != null && supportedCodecs.contains(codec)) {
+ return CompressionCodec.Factory.INSTANCE.createCodec(codec);
+ }
+ }
+ return null;
+ }
+
+ private static CompressionUtil.CodecType parseCodec(String name) {
+ switch (name.trim().toLowerCase(Locale.ROOT)) {
+ case "lz4":
+ case "lz4_frame":
+ return CompressionUtil.CodecType.LZ4_FRAME;
+ case "zstd":
+ return CompressionUtil.CodecType.ZSTD;
+ default:
+ return null;
+ }
+ }
+}
diff --git a/flight/flight-core/src/main/java/org/apache/arrow/flight/OutboundStreamListenerImpl.java b/flight/flight-core/src/main/java/org/apache/arrow/flight/OutboundStreamListenerImpl.java
index a1bde3a848..c440abb021 100644
--- a/flight/flight-core/src/main/java/org/apache/arrow/flight/OutboundStreamListenerImpl.java
+++ b/flight/flight-core/src/main/java/org/apache/arrow/flight/OutboundStreamListenerImpl.java
@@ -22,6 +22,7 @@
import org.apache.arrow.util.Preconditions;
import org.apache.arrow.vector.VectorSchemaRoot;
import org.apache.arrow.vector.VectorUnloader;
+import org.apache.arrow.vector.compression.CompressionCodec;
import org.apache.arrow.vector.dictionary.DictionaryProvider;
import org.apache.arrow.vector.ipc.message.IpcOption;
@@ -32,12 +33,21 @@ abstract class OutboundStreamListenerImpl implements OutboundStreamListener {
protected volatile VectorUnloader unloader; // null until stream started
protected IpcOption option; // null until stream started
protected boolean tryZeroCopy = ArrowMessage.ENABLE_ZERO_COPY_WRITE;
+ private final CompressionCodec compressionCodec;
OutboundStreamListenerImpl(
FlightDescriptor descriptor, CallStreamObserver responseObserver) {
+ this(descriptor, responseObserver, null);
+ }
+
+ OutboundStreamListenerImpl(
+ FlightDescriptor descriptor,
+ CallStreamObserver responseObserver,
+ CompressionCodec compressionCodec) {
Preconditions.checkNotNull(responseObserver, "responseObserver must be provided");
this.descriptor = descriptor;
this.responseObserver = responseObserver;
+ this.compressionCodec = compressionCodec;
this.unloader = null;
}
@@ -56,7 +66,12 @@ public void start(VectorSchemaRoot root, DictionaryProvider dictionaries, IpcOpt
this.option = option;
try {
DictionaryUtils.generateSchemaMessages(
- root.getSchema(), descriptor, dictionaries, option, responseObserver::onNext);
+ root.getSchema(),
+ descriptor,
+ dictionaries,
+ option,
+ compressionCodec,
+ responseObserver::onNext);
} catch (RuntimeException e) {
// Propagate runtime exceptions, like those raised when trying to write unions with V4
// metadata
@@ -68,7 +83,9 @@ public void start(VectorSchemaRoot root, DictionaryProvider dictionaries, IpcOpt
throw new RuntimeException("Could not generate and send all schema messages", e);
}
// We include the null count and align buffers to be compatible with Flight/C++
- unloader = new VectorUnloader(root, /* includeNullCount */ true, /* alignBuffers */ true);
+ unloader =
+ new VectorUnloader(
+ root, /* includeNullCount */ true, compressionCodec, /* alignBuffers */ true);
}
@Override
diff --git a/flight/flight-core/src/test/java/org/apache/arrow/flight/TestIpcCompression.java b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestIpcCompression.java
new file mode 100644
index 0000000000..42002eae83
--- /dev/null
+++ b/flight/flight-core/src/test/java/org/apache/arrow/flight/TestIpcCompression.java
@@ -0,0 +1,93 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.arrow.flight;
+
+import static org.apache.arrow.flight.FlightTestUtil.LOCALHOST;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.nio.charset.StandardCharsets;
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.apache.arrow.vector.VarCharVector;
+import org.apache.arrow.vector.VectorSchemaRoot;
+import org.apache.arrow.vector.compression.CompressionUtil;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
+
+public class TestIpcCompression {
+ private static final byte[] VALUE = "compressible-value".getBytes(StandardCharsets.UTF_8);
+
+ @ParameterizedTest
+ @EnumSource(
+ value = CompressionUtil.CodecType.class,
+ names = {"LZ4_FRAME", "ZSTD"})
+ public void negotiatesCompression(CompressionUtil.CodecType codec) throws Exception {
+ assertRoundTrip(codec, true, codec);
+ }
+
+ @ParameterizedTest
+ @EnumSource(
+ value = CompressionUtil.CodecType.class,
+ names = {"LZ4_FRAME", "ZSTD"})
+ public void fallsBackForServersWithoutCompression(CompressionUtil.CodecType codec)
+ throws Exception {
+ assertRoundTrip(codec, false, CompressionUtil.CodecType.NO_COMPRESSION);
+ }
+
+ private static void assertRoundTrip(
+ CompressionUtil.CodecType requested,
+ boolean enableServerCompression,
+ CompressionUtil.CodecType expected)
+ throws Exception {
+ try (BufferAllocator allocator = new RootAllocator(Long.MAX_VALUE)) {
+ final FlightServer.Builder serverBuilder =
+ FlightServer.builder(
+ allocator,
+ Location.forGrpcInsecure(LOCALHOST, 0),
+ new NoOpFlightProducer() {
+ @Override
+ public void getStream(
+ CallContext context, Ticket ticket, ServerStreamListener listener) {
+ try (VarCharVector vector = new VarCharVector("value", allocator);
+ VectorSchemaRoot root = VectorSchemaRoot.of(vector)) {
+ vector.allocateNew();
+ vector.setSafe(0, VALUE);
+ vector.setValueCount(1);
+ root.setRowCount(1);
+ listener.start(root);
+ listener.putNext();
+ listener.completed();
+ }
+ }
+ });
+ if (enableServerCompression) {
+ serverBuilder.ipcCompression(
+ CompressionUtil.CodecType.LZ4_FRAME, CompressionUtil.CodecType.ZSTD);
+ }
+ try (FlightServer server = serverBuilder.build().start();
+ FlightClient client = FlightClient.builder(allocator, server.getLocation()).build();
+ FlightStream stream =
+ client.getStream(
+ new Ticket(new byte[0]), CallOptions.acceptIpcCompression(requested))) {
+ assertTrue(stream.next());
+ assertEquals("compressible-value", stream.getRoot().getVector(0).getObject(0).toString());
+ assertEquals(expected, stream.compressionType);
+ }
+ }
+ }
+}
diff --git a/flight/flight-sql-jdbc-core/pom.xml b/flight/flight-sql-jdbc-core/pom.xml
index d6fa11688d..f5ba51ae20 100644
--- a/flight/flight-sql-jdbc-core/pom.xml
+++ b/flight/flight-sql-jdbc-core/pom.xml
@@ -79,6 +79,12 @@ under the License.
${arrow.vector.classifier}
+