Java
Maven dependency
<dependencies>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-netty-shaded</artifactId>
<version>1.65.0</version>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-protobuf</artifactId>
<version>1.65.0</version>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-stub</artifactId>
<version>1.65.0</version>
</dependency>
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>4.27.0</version>
</dependency>
</dependencies>
Add the protoc Maven plugin to generate stubs from the proto files during build, or generate them manually and commit to your repo.
Connect
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import com.tinybroker.v1.BrokerGrpc;
ManagedChannel channel = ManagedChannelBuilder
.forAddress("tinybroker", 50051)
.usePlaintext()
.build();
// Blocking stub for unary RPCs
BrokerGrpc.BrokerBlockingStub blockingStub = BrokerGrpc.newBlockingStub(channel);
// Async stub for streaming RPCs
BrokerGrpc.BrokerStub asyncStub = BrokerGrpc.newStub(channel);
Publish
import com.tinybroker.v1.ClientProto.PublishRequest;
import com.google.protobuf.ByteString;
blockingStub.publish(PublishRequest.newBuilder()
.setTopic("events.user.signup")
.setPayload(ByteString.copyFromUtf8("{\"user_id\":\"abc123\"}"))
.build());
Subscribe (fan-out)
import com.tinybroker.v1.ClientProto.*;
import com.tinybroker.v1.ServerProto.SubscriptionEvent;
import com.tinybroker.v1.SharedProto.DeliveryMode;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.CountDownLatch;
CountDownLatch latch = new CountDownLatch(1);
StreamObserver<SubscriptionEvent> responseObserver = new StreamObserver<>() {
@Override
public void onNext(SubscriptionEvent event) {
switch (event.getEventCase()) {
case SUBSCRIPTION_ID:
System.out.println("confirmed: " + event.getSubscriptionId());
break;
case MESSAGE:
var msg = event.getMessage();
System.out.printf("topic=%s payload=%s%n",
msg.getTopic(), msg.getPayload().toStringUtf8());
break;
case PATTERNS_UPDATED:
System.out.println("patterns: " + event.getPatternsUpdated().getTopicPatternsList());
break;
}
}
@Override
public void onError(Throwable t) { latch.countDown(); /* reconnect */ }
@Override
public void onCompleted() { latch.countDown(); }
};
StreamObserver<SubscribeCommand> requestObserver =
asyncStub.subscribe(responseObserver);
requestObserver.onNext(SubscribeCommand.newBuilder()
.setOpen(OpenSubscription.newBuilder()
.addTopicPatterns("events.user.*")
.addTopicPatterns("events.order.*")
.setMode(DeliveryMode.DELIVERY_MODE_MULTI)
.build())
.build());
latch.await(); // block until stream ends
Subscribe (consumer group)
StreamObserver<SubscribeCommand> requestObserver =
asyncStub.subscribe(new StreamObserver<>() {
String subscriptionId;
@Override
public void onNext(SubscriptionEvent event) {
if (event.hasSubscriptionId()) {
subscriptionId = event.getSubscriptionId();
} else if (event.hasMessage()) {
var msg = event.getMessage();
processJob(msg.getPayload().toByteArray());
blockingStub.ack(AckRequest.newBuilder()
.setSubscriptionId(subscriptionId)
.setSystemId(msg.getSystemId())
.build());
}
}
@Override public void onError(Throwable t) { /* reconnect */ }
@Override public void onCompleted() {}
});
requestObserver.onNext(SubscribeCommand.newBuilder()
.setOpen(OpenSubscription.newBuilder()
.setSubscriptionId("job-workers")
.addTopicPatterns("jobs.*")
.setMode(DeliveryMode.DELIVERY_MODE_BALANCED)
.setConsumerGroup("job-workers")
.build())
.build());
Dynamic pattern management
// Add patterns
requestObserver.onNext(SubscribeCommand.newBuilder()
.setAddPatterns(AddTopicPatterns.newBuilder()
.addTopicPatterns("alerts.*")
.build())
.build());
// Remove patterns
requestObserver.onNext(SubscribeCommand.newBuilder()
.setRemovePatterns(RemoveTopicPatterns.newBuilder()
.addTopicPatterns("events.user.*")
.build())
.build());
Shutdown
channel.shutdown().awaitTermination(5, TimeUnit.SECONDS);