# Google Cloud Pub/Sub (/connectors/gcp-pub-sub)



Point your Google Cloud Pub/Sub application at KubeMQ by setting **one environment
variable** — `PUBSUB_EMULATOR_HOST`. The **Google Cloud Pub/Sub connector** is a built-in,
wire-protocol bridge inside kubemq-server that speaks the genuine Pub/Sub v1 gRPC services on
a dedicated gRPC listener (default port **8085**, the Pub/Sub emulator convention). Any
standard, unmodified Pub/Sub client — the Go, Python, Java, Node.js, C#, and Ruby first-party
Google clients, plus `gcloud pubsub` — talks to KubeMQ exactly as it would to Google's local
emulator, with **no code changes, no library swap, and no KubeMQ SDK**.

## What is the Pub/Sub connector [#what-is-the-pubsub-connector]

The connector is a **single gRPC listener** that implements the real Pub/Sub v1 wire
protocol — **38 RPCs** across four services:

* **`google.pubsub.v1.Publisher`** (9 RPCs) — topics and publish.
* **`google.pubsub.v1.Subscriber`** (16 RPCs) — subscriptions, pull, streaming pull, ack,
  snapshots, and seek.
* **`google.pubsub.v1.SchemaService`** (10 RPCs) — Avro and Protobuf schemas.
* **`google.iam.v1.IAMPolicy`** (3 permissive stubs) — emulator-parity IAM.

Every official Pub/Sub client library honours the `PUBSUB_EMULATOR_HOST` environment
variable: when it is set, the SDK **clears credentials, skips Google auth, and dials insecure
gRPC** — exactly as it would against Google's local emulator. **The connector *is* the
emulator** — there is no separate emulator to install, no LocalStack, and no boot-the-server
step beyond running kubemq-server.

Two KubeMQ primitives back the model: a **topic** maps onto a native KubeMQ **Events Store**
log `gcp.{topic}` (the authoritative, replayable source of truth), and each **subscription**
maps onto a native **Queue** channel `gcp.sub.{subscription}`. A publish is written **once**
to the topic log, then fanned out to one queue copy per subscription — so Pub/Sub producers
and native gRPC/REST consumers interoperate on the same messages.

<Callout type="warn">
  The Pub/Sub connector is &#x2A;*opt-in (disabled by default)**. A stock kubemq-server does **not**
  bind gRPC port 8085 until you enable it with `CONNECTORS_GCP_ENABLE=true` (Docker) or
  `spec.gcp.enabled: true` (Kubernetes). See
  [Getting started](/connectors/gcp-pub-sub/tutorials/getting-started).
</Callout>

## How it works [#how-it-works]

A Pub/Sub SDK client dials the connector's gRPC endpoint. A `Publish` is written once to the
topic's Events Store log `gcp.{topic}` through the message broker, then fanned out to one
Queue copy per subscription `gcp.sub.{subscription}` (applying each subscription's filter).
`Pull` / `StreamingPull` lease each delivered message under an ack-deadline, and `Acknowledge`
removes it from the subscription queue.

<Mermaid
  chart="`
graph LR
APP[&#x22;Pub/Sub SDK client<br/>PUBSUB_EMULATOR_HOST=:8085&#x22;]
CONN[&#x22;Pub/Sub connector<br/>:8085 (gRPC)&#x22;]
LOG[(&#x22;Events Store log<br/>gcp.orders&#x22;)]
Q{{&#x22;Queue channel<br/>gcp.sub.sub-orders&#x22;}}
BROKER[&#x22;Message Broker&#x22;]
SUB[&#x22;Consumer<br/>(Pull / gRPC / REST)&#x22;]

APP -- &#x22;Publish&#x22; --> CONN
CONN -- &#x22;SendEventsStore&#x22; --> LOG
LOG -- &#x22;fan-out per subscription&#x22; --> Q
CONN -- &#x22;Pull / StreamingPull&#x22; --> Q
LOG --> BROKER
Q --> BROKER
BROKER -. &#x22;deliver + ack_id&#x22; .-> SUB

class APP,SUB client
class CONN connector
class LOG,Q,BROKER broker
`"
/>

*A publish writes once to the Events Store log `gcp.{topic}` and fans out one Queue copy per subscription on `gcp.sub.{subscription}`, all backed by the message broker; consumers pull each message under an ack-deadline lease.*

## Ports & protocol surface [#ports--protocol-surface]

| Port   | Transport     | Protocol                        | Notes                                                                                                                                                                                                            |
| ------ | ------------- | ------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `8085` | Insecure gRPC | Pub/Sub v1 gRPC (emulator mode) | The Pub/Sub emulator convention. **Opt-in** — bound only when `CONNECTORS_GCP_ENABLE=true`. Must differ from the gRPC/REST/HTTP and AWS-connector ports. No auth, no TLS — exactly like Google's local emulator. |
| —      | gRPC over TLS | Pub/Sub v1 gRPC                 | TLS is provided by the server-wide `Security` block — there is **no Pub/Sub-specific TLS option**. The default emulator path is unencrypted; do not expose port 8085 to untrusted networks.                      |

The connector is **gRPC only** — there is **no REST/JSON v1** surface (no grpc-gateway).
Clients and tools that only speak the Pub/Sub REST API will not work; use a gRPC client
library or `gcloud` (which uses gRPC for the emulator). See
[Architecture](/connectors/gcp-pub-sub/concepts/architecture) for the dispatch detail.

## Send a message [#send-a-message]

The example below runs the full round-trip — `CreateTopic` → `CreateSubscription` →
`Publish` → `Pull` → `Acknowledge` — over a stock Google Cloud Pub/Sub client. The only
change versus a real-GCP app is the **emulator endpoint**: each client reads
`PUBSUB_EMULATOR_HOST` (default `localhost:8085`) and a project id from `PUBSUB_PROJECT_ID`
(any value — the project segment is parsed but ignored).

<Tabs groupId="language" items="['Go','Python','Java','JavaScript','C#','Ruby']">
  <Tab value="Go">
    ```go
    package main

    import (
    	"context"
    	"fmt"
    	"log"
    	"os"
    	"sync"
    	"time"

    	"cloud.google.com/go/pubsub"
    )

    func projectID() string {
    	if v := os.Getenv("PUBSUB_PROJECT_ID"); v != "" {
    		return v
    	}
    	return "my-project" // any id; the project segment is parsed but ignored.
    }

    func main() {
    	ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
    	defer cancel()

    	// PUBSUB_EMULATOR_HOST (default localhost:8085) routes the official client at
    	// the connector over insecure gRPC with credentials cleared — no code change.
    	client, err := pubsub.NewClient(ctx, projectID())
    	if err != nil {
    		log.Fatalf("NewClient: %v", err)
    	}
    	defer client.Close()

    	// CreateTopic "orders" maps to Events Store log "gcp.orders".
    	topic, err := client.CreateTopic(ctx, "orders")
    	if err != nil {
    		log.Fatalf("CreateTopic: %v", err)
    	}
    	defer topic.Stop()

    	// CreateSubscription maps to Queue channel "gcp.sub.sub-orders".
    	sub, err := client.CreateSubscription(ctx, "sub-orders", pubsub.SubscriptionConfig{
    		Topic:       topic,
    		AckDeadline: 10 * time.Second, // connector default; valid range 10..600s.
    	})
    	if err != nil {
    		log.Fatalf("CreateSubscription: %v", err)
    	}

    	id, err := topic.Publish(ctx, &pubsub.Message{Data: []byte("hello from cloud.google.com/go/pubsub")}).Get(ctx)
    	if err != nil {
    		log.Fatalf("Publish: %v", err)
    	}
    	fmt.Printf("published: %s\n", id)

    	recvCtx, recvCancel := context.WithTimeout(ctx, 15*time.Second)
    	defer recvCancel()
    	var once sync.Once
    	err = sub.Receive(recvCtx, func(_ context.Context, m *pubsub.Message) {
    		fmt.Printf("received: %q\n", string(m.Data))
    		m.Ack() // Acknowledge by ack_id under its lease.
    		once.Do(recvCancel)
    	})
    	if err != nil && recvCtx.Err() == nil {
    		log.Fatalf("Receive: %v", err)
    	}
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python
    import os

    from google.cloud import pubsub_v1


    def project_id() -> str:
        # Any id works — the project segment is parsed but ignored.
        return os.environ.get("PUBSUB_PROJECT_ID", "my-project")


    def main() -> None:
        # Both clients honour PUBSUB_EMULATOR_HOST (default localhost:8085): they clear
        # credentials, skip Google auth, and dial insecure gRPC.
        publisher = pubsub_v1.PublisherClient()
        subscriber = pubsub_v1.SubscriberClient()

        proj = project_id()
        topic_path = publisher.topic_path(proj, "orders")          # -> gcp.orders
        sub_path = subscriber.subscription_path(proj, "sub-orders")  # -> gcp.sub.sub-orders

        publisher.create_topic(request={"name": topic_path})
        subscriber.create_subscription(request={"name": sub_path, "topic": topic_path})

        future = publisher.publish(topic_path, b"hello from google-cloud-pubsub")
        print(f"published: {future.result(timeout=15)}")

        resp = subscriber.pull(request={"subscription": sub_path, "max_messages": 1}, timeout=20)
        msg = resp.received_messages[0]
        print(f"received: {msg.message.data.decode()!r}")

        # Acknowledge by ack_id; the message leaves the subscription queue.
        subscriber.acknowledge(request={"subscription": sub_path, "ack_ids": [msg.ack_id]})
        subscriber.close()


    if __name__ == "__main__":
        main()
    ```
  </Tab>

  <Tab value="Java">
    ```java
    import com.google.api.gax.core.NoCredentialsProvider;
    import com.google.api.gax.grpc.GrpcTransportChannel;
    import com.google.api.gax.rpc.FixedTransportChannelProvider;
    import com.google.api.gax.rpc.TransportChannelProvider;
    import com.google.cloud.pubsub.v1.SubscriptionAdminClient;
    import com.google.cloud.pubsub.v1.SubscriptionAdminSettings;
    import com.google.cloud.pubsub.v1.TopicAdminClient;
    import com.google.cloud.pubsub.v1.TopicAdminSettings;
    import com.google.cloud.pubsub.v1.stub.GrpcSubscriberStub;
    import com.google.cloud.pubsub.v1.stub.SubscriberStubSettings;
    import com.google.protobuf.ByteString;
    import com.google.pubsub.v1.AcknowledgeRequest;
    import com.google.pubsub.v1.PublishRequest;
    import com.google.pubsub.v1.PubsubMessage;
    import com.google.pubsub.v1.PullRequest;
    import com.google.pubsub.v1.PullResponse;
    import com.google.pubsub.v1.PushConfig;
    import com.google.pubsub.v1.ReceivedMessage;
    import com.google.pubsub.v1.SubscriptionName;
    import com.google.pubsub.v1.TopicName;
    import io.grpc.ManagedChannel;
    import io.grpc.ManagedChannelBuilder;

    public final class Main {
        public static void main(String[] args) throws Exception {
            String emulatorHost = System.getenv().getOrDefault("PUBSUB_EMULATOR_HOST", "localhost:8085");
            String projectId = System.getenv().getOrDefault("PUBSUB_PROJECT_ID", "my-project");

            // The Go/Python/Node clients auto-detect the emulator; Java points a plaintext
            // gRPC channel at the host explicitly with NoCredentialsProvider.
            ManagedChannel channel = ManagedChannelBuilder.forTarget(emulatorHost).usePlaintext().build();
            TransportChannelProvider channelProvider =
                    FixedTransportChannelProvider.create(GrpcTransportChannel.create(channel));
            NoCredentialsProvider noCreds = NoCredentialsProvider.create();

            TopicName topic = TopicName.of(projectId, "orders");                 // -> gcp.orders
            SubscriptionName sub = SubscriptionName.of(projectId, "sub-orders"); // -> gcp.sub.sub-orders

            try (TopicAdminClient topicAdmin = TopicAdminClient.create(TopicAdminSettings.newBuilder()
                    .setTransportChannelProvider(channelProvider).setCredentialsProvider(noCreds).build());
                    SubscriptionAdminClient subAdmin = SubscriptionAdminClient.create(
                            SubscriptionAdminSettings.newBuilder()
                                    .setTransportChannelProvider(channelProvider)
                                    .setCredentialsProvider(noCreds).build());
                    GrpcSubscriberStub subStub = GrpcSubscriberStub.create(SubscriberStubSettings.newBuilder()
                            .setTransportChannelProvider(channelProvider)
                            .setCredentialsProvider(noCreds).build())) {

                topicAdmin.createTopic(topic);
                subAdmin.createSubscription(sub, topic, PushConfig.getDefaultInstance(), 10);

                topicAdmin.publish(PublishRequest.newBuilder()
                        .setTopic(topic.toString())
                        .addMessages(PubsubMessage.newBuilder()
                                .setData(ByteString.copyFromUtf8("hello from google-cloud-pubsub")).build())
                        .build());

                PullResponse resp = subStub.pullCallable().call(PullRequest.newBuilder()
                        .setSubscription(sub.toString()).setMaxMessages(1).build());
                ReceivedMessage got = resp.getReceivedMessages(0);
                System.out.printf("received: %s%n", got.getMessage().getData().toStringUtf8());

                // Acknowledge by ack_id.
                subStub.acknowledgeCallable().call(AcknowledgeRequest.newBuilder()
                        .setSubscription(sub.toString()).addAckIds(got.getAckId()).build());
            } finally {
                channel.shutdown();
            }
        }
    }
    ```
  </Tab>

  <Tab value="JavaScript">
    ```typescript
    import { PubSub, v1 } from "@google-cloud/pubsub";

    const projectId = process.env["PUBSUB_PROJECT_ID"] ?? "my-project";

    // The high-level PubSub client reads PUBSUB_EMULATOR_HOST (default localhost:8085)
    // and resolves the insecure emulator transport; reuse its options for the v1 clients.
    const baseOptions = new PubSub({ projectId }).options;
    const options = { ...baseOptions, port: baseOptions.port === undefined ? undefined : Number(baseOptions.port) };
    const publisher = new v1.PublisherClient(options);
    const subscriber = new v1.SubscriberClient(options);

    async function main(): Promise<void> {
      const topic = publisher.projectTopicsPath(projectId, "orders");          // -> gcp.orders
      const sub = subscriber.subscriptionPath(projectId, "sub-orders");        // -> gcp.sub.sub-orders

      await publisher.createTopic({ name: topic });
      await subscriber.createSubscription({ name: sub, topic, ackDeadlineSeconds: 10 });

      const [published] = await publisher.publish({
        topic,
        messages: [{ data: Buffer.from("hello from @google-cloud/pubsub") }],
      });
      console.log(`published: ${published.messageIds?.[0]}`);

      const [pull] = await subscriber.pull({ subscription: sub, maxMessages: 1 });
      const received = pull.receivedMessages![0];
      console.log(`received: ${Buffer.from(received.message!.data as Uint8Array).toString("utf8")}`);

      // Acknowledge by ack_id.
      await subscriber.acknowledge({ subscription: sub, ackIds: [received.ackId!] });
    }

    main().catch((err) => {
      console.error(err);
      process.exit(1);
    });
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    using Google.Api.Gax;
    using Google.Cloud.PubSub.V1;
    using Google.Protobuf;

    var projectId = Environment.GetEnvironmentVariable("PUBSUB_PROJECT_ID") ?? "my-project";
    var topicName = TopicName.FromProjectTopic(projectId, "orders");           // -> gcp.orders
    var subName = SubscriptionName.FromProjectSubscription(projectId, "sub-orders"); // -> gcp.sub.sub-orders

    // The .NET client does NOT auto-detect the emulator — set EmulatorOnly so each
    // client reads PUBSUB_EMULATOR_HOST (default localhost:8085) and dials insecurely.
    var publisher = await new PublisherServiceApiClientBuilder
    {
        EmulatorDetection = EmulatorDetection.EmulatorOnly,
    }.BuildAsync();
    var subscriber = await new SubscriberServiceApiClientBuilder
    {
        EmulatorDetection = EmulatorDetection.EmulatorOnly,
    }.BuildAsync();

    await publisher.CreateTopicAsync(topicName);
    await subscriber.CreateSubscriptionAsync(subName, topicName, pushConfig: null, ackDeadlineSeconds: 10);

    var publishResponse = await publisher.PublishAsync(topicName, new[]
    {
        new PubsubMessage { Data = ByteString.CopyFromUtf8("hello from Google.Cloud.PubSub.V1") },
    });
    Console.WriteLine($"published: {publishResponse.MessageIds[0]}");

    var pull = await subscriber.PullAsync(subName, maxMessages: 1);
    var received = pull.ReceivedMessages[0];
    Console.WriteLine($"received: {received.Message.Data.ToStringUtf8()}");

    // Acknowledge by ack_id.
    await subscriber.AcknowledgeAsync(subName, new[] { received.AckId });
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby
    # frozen_string_literal: true

    require "google/cloud/pubsub"

    project_id = ENV["PUBSUB_PROJECT_ID"] || "my-project"
    # Ruby needs the emulator host passed explicitly (it does not always read the env var).
    emulator_host = ENV["PUBSUB_EMULATOR_HOST"] || "localhost:8085"

    pubsub = Google::Cloud::PubSub.new(project_id: project_id, emulator_host: emulator_host)
    topic_admin = pubsub.topic_admin
    sub_admin   = pubsub.subscription_admin

    topic_path = pubsub.topic_path("orders")            # -> gcp.orders
    sub_path   = pubsub.subscription_path("sub-orders") # -> gcp.sub.sub-orders

    topic = topic_admin.create_topic(name: topic_path)
    sub_admin.create_subscription(name: sub_path, topic: topic_path, ack_deadline_seconds: 10)

    publisher = pubsub.publisher(topic.name)
    msg = publisher.publish("hello from google-cloud-pubsub")
    puts "published: #{msg.message_id}"

    subscriber = pubsub.subscriber(sub_path)
    received = subscriber.pull(immediate: false, max: 1)
    rcv = received.first
    puts "received: #{rcv.data.inspect}"

    # Acknowledge by ack_id; the message leaves the subscription queue.
    rcv.acknowledge!
    ```
  </Tab>
</Tabs>

## Supported languages [#supported-languages]

The connector speaks the genuine Pub/Sub v1 wire protocol, so any first-party Google Cloud
Pub/Sub client works — you only set `PUBSUB_EMULATOR_HOST`. There is no KubeMQ SDK, no proto
bindings, and no published package; the examples pin one official Google client per language.

| Language                | Google Cloud Pub/Sub client                                                                                                      | Emulator construction                                                                      |
| ----------------------- | -------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------ |
| Go                      | [`cloud.google.com/go/pubsub`](https://pkg.go.dev/cloud.google.com/go/pubsub)                                                    | `pubsub.NewClient(ctx, projectID)` — auto-detects `PUBSUB_EMULATOR_HOST`, dials insecurely |
| Python                  | [`google-cloud-pubsub`](https://cloud.google.com/python/docs/reference/pubsub/latest) (via **uv**)                               | `pubsub_v1.PublisherClient()` — honours `PUBSUB_EMULATOR_HOST`                             |
| Java                    | [`com.google.cloud:google-cloud-pubsub`](https://cloud.google.com/java/docs/reference/google-cloud-pubsub/latest/overview) (BOM) | plaintext `ManagedChannel` to the emulator host + `NoCredentialsProvider`                  |
| JavaScript / TypeScript | [`@google-cloud/pubsub`](https://cloud.google.com/nodejs/docs/reference/pubsub/latest)                                           | `new PubSub({ projectId })` — auto-detects the emulator; run via `tsx`                     |
| C# / .NET               | [`Google.Cloud.PubSub.V1`](https://cloud.google.com/dotnet/docs/reference/Google.Cloud.PubSub.V1/latest) (.NET 8)                | `…Builder { EmulatorDetection = EmulatorDetection.EmulatorOnly }.Build()`                  |
| Ruby                    | [`google-cloud-pubsub`](https://cloud.google.com/ruby/docs/reference/google-cloud-pubsub/latest)                                 | `Google::Cloud::PubSub.new(project_id:, emulator_host:)`                                   |

<Callout type="info">
  There is **no Rust tab** — Google ships no first-party Pub/Sub client for Rust. The connector's
  verified example suite is the six languages above. Most clients auto-detect the emulator from
  `PUBSUB_EMULATOR_HOST`; &#x2A;*C#** needs `EmulatorDetection.EmulatorOnly`, **Ruby** needs an
  explicit `emulator_host:`, and **Java** points a plaintext channel at the host — see
  [Connectivity & emulator mode](/connectors/gcp-pub-sub/how-to/connectivity-and-emulator-mode).
</Callout>

## Next steps [#next-steps]

<Cards>
  <Card title="Getting started" href="/connectors/gcp-pub-sub/tutorials/getting-started" description="Point your Pub/Sub SDK at port 8085 via PUBSUB_EMULATOR_HOST and run a publish-and-pull round-trip in minutes." />

  <Card title="Configuration" href="/connectors/gcp-pub-sub/concepts/configuration" description="The opt-in CONNECTORS_GCP_ENABLE variable, gRPC port 8085, and the thirteen CONNECTORS_GCP_* settings." />

  <Card title="Publish & Subscribe" href="/connectors/gcp-pub-sub/how-to/publish-subscribe" description="The core topic-to-subscription round-trip mapped onto an Events Store log and a Queue channel." />

  <Card title="Fan-Out" href="/connectors/gcp-pub-sub/how-to/fan-out" description="One publish, many subscriptions — fan a single topic message out to multiple independent subscriptions." />
</Cards>
