Skip to content

Commit d870839

Browse files
committed
feat: Allow REST transport for PubSub Java client.
This leverages capabilities added in googleapis#1162
1 parent 94c55d0 commit d870839

1 file changed

Lines changed: 23 additions & 7 deletions

File tree

  • google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1

google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java

Lines changed: 23 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -36,13 +36,17 @@
3636
import com.google.api.gax.core.FixedExecutorProvider;
3737
import com.google.api.gax.core.InstantiatingExecutorProvider;
3838
import com.google.api.gax.grpc.GrpcCallContext;
39+
import com.google.api.gax.httpjson.HttpJsonCallContext;
40+
import com.google.api.gax.httpjson.InstantiatingHttpJsonChannelProvider;
3941
import com.google.api.gax.retrying.RetrySettings;
42+
import com.google.api.gax.rpc.ApiCallContext;
4043
import com.google.api.gax.rpc.HeaderProvider;
4144
import com.google.api.gax.rpc.NoHeaderProvider;
4245
import com.google.api.gax.rpc.StatusCode;
4346
import com.google.api.gax.rpc.TransportChannelProvider;
4447
import com.google.auth.oauth2.GoogleCredentials;
4548
import com.google.cloud.pubsub.v1.stub.GrpcPublisherStub;
49+
import com.google.cloud.pubsub.v1.stub.HttpJsonPublisherStub;
4650
import com.google.cloud.pubsub.v1.stub.PublisherStub;
4751
import com.google.cloud.pubsub.v1.stub.PublisherStubSettings;
4852
import com.google.common.base.Preconditions;
@@ -120,9 +124,10 @@ public class Publisher implements PublisherInterface {
120124

121125
private final boolean enableCompression;
122126
private final long compressionBytesThreshold;
127+
private final boolean enableRESTJsonTransport;
123128

124-
private final GrpcCallContext publishContext;
125-
private final GrpcCallContext publishContextWithCompression;
129+
private final ApiCallContext publishContext;
130+
private final ApiCallContext publishContextWithCompression;
126131

127132
/** The maximum number of messages in one request. Defined by the API. */
128133
public static long getApiMaxRequestElementCount() {
@@ -152,6 +157,8 @@ private Publisher(Builder builder) throws IOException {
152157
this.messageTransform = builder.messageTransform;
153158
this.enableCompression = builder.enableCompression;
154159
this.compressionBytesThreshold = builder.compressionBytesThreshold;
160+
this.enableRESTJsonTransport =
161+
builder.channelProvider instanceof InstantiatingHttpJsonChannelProvider;
155162

156163
messagesBatches = new HashMap<>();
157164
messagesBatchLock = new ReentrantLock();
@@ -199,15 +206,24 @@ private Publisher(Builder builder) throws IOException {
199206
StatusCode.Code.UNAVAILABLE)
200207
.setRetrySettings(retrySettingsBuilder.build())
201208
.setBatchingSettings(BatchingSettings.newBuilder().setIsEnabled(false).build());
202-
this.publisherStub = GrpcPublisherStub.create(stubSettings.build());
209+
this.publisherStub =
210+
this.enableRESTJsonTransport
211+
? HttpJsonPublisherStub.create(stubSettings.build())
212+
: GrpcPublisherStub.create(stubSettings.build());
203213
backgroundResourceList.add(publisherStub);
204214
backgroundResources = new BackgroundResourceAggregation(backgroundResourceList);
205215
shutdown = new AtomicBoolean(false);
206216
messagesWaiter = new Waiter();
207-
this.publishContext = GrpcCallContext.createDefault();
217+
this.publishContext =
218+
this.enableRESTJsonTransport
219+
? HttpJsonCallContext.createDefault()
220+
: GrpcCallContext.createDefault();
208221
this.publishContextWithCompression =
209-
GrpcCallContext.createDefault()
210-
.withCallOptions(CallOptions.DEFAULT.withCompression(GZIP_COMPRESSION));
222+
this.enableRESTJsonTransport
223+
? this.publishContext
224+
: // TODO
225+
GrpcCallContext.createDefault()
226+
.withCallOptions(CallOptions.DEFAULT.withCompression(GZIP_COMPRESSION));
211227
}
212228

213229
/** Topic which the publisher publishes to. */
@@ -448,7 +464,7 @@ private void publishAllWithoutInflightForKey(final String orderingKey) {
448464
}
449465

450466
private ApiFuture<PublishResponse> publishCall(OutstandingBatch outstandingBatch) {
451-
GrpcCallContext context = publishContext;
467+
ApiCallContext context = publishContext;
452468
if (enableCompression && outstandingBatch.batchSizeBytes >= compressionBytesThreshold) {
453469
context = publishContextWithCompression;
454470
}

0 commit comments

Comments
 (0)