Single-node · In-memory · gRPC

Heads up: I ended up switching to NATS + PostgreSQL. If you need consumer groups, bare NATS won't help — you need JetStream, which requires persistence. tinybroker fills the gap when you want in-memory consumer groups with no persistence overhead. The code works and is free to use; just know why you're reaching for it.

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);