Migrating from Google Cloud Pub/Sub
Set PUBSUB_EMULATOR_HOST to KubeMQ: quick start, topic and subscription mapping, KubeMQ token auth, no Google IAM or TLS, node-local ordering.
Point your existing Google Cloud Pub/Sub application at KubeMQ by changing one environment
variable. The KubeMQ GCP connector exposes an emulator-compatible Pub/Sub v1 gRPC surface, so
unmodified google-cloud-pubsub 2.x clients, gcloud pubsub, and any SDK that honours
PUBSUB_EMULATOR_HOST connect with no application code change — set the env var, clear the
Google credentials, and you are on KubeMQ. There is no KubeMQ SDK to adopt, no proto, and no data
migration: topics live in normal KubeMQ Events Store logs and subscriptions in normal Queue
channels. Rollback is config-only (CONNECTORS_GCP_ENABLE=false).
This is an endpoint-only drop-in. But several connector behaviors deviate from real Google Pub/Sub — read What Does NOT Migrate / Deviations before you cut over; most are invisible until a corner case hits production.
Overview
KubeMQ's GCP connector exposes an emulator-compatible gRPC surface that accepts unmodified
google-cloud-pubsub 2.x client libraries, gcloud pubsub, and any other SDK that honours
PUBSUB_EMULATOR_HOST. No application code changes are required — point the env var at KubeMQ and
clear the Google credentials.
| Attribute | Value |
|---|---|
| Protocol | gRPC (Pub/Sub v1 — Publisher + Subscriber + SchemaService + IAMPolicy) |
| Default port | 8085 (emulator convention) |
| Canonical client | google-cloud-pubsub 2.x (pubsub_v1 gRPC stubs) |
| Drop-in level | endpoint-only — set PUBSUB_EMULATOR_HOST, clear credentials |
| Enable | CONNECTORS_GCP_ENABLE=true; Kubernetes spec.gcp.enabled: true |
Quick start: repoint one client
Try the switch on your own machine before you plan the migration.
The Docker tab points one client at the KubeMQ you started with Try KubeMQ, turning its connector on where needed and keeping your evaluation or license and your data. The Kubernetes tab works on a cluster you installed with Install on Kubernetes. Use it to try the switch, not for production; see Plans compared for staging and production licenses. If your current broker runs on this machine, stop it first, because two programs cannot publish the same port. The kmq command runs natively on Windows; the other commands are bash, so on Windows run them in WSL (Windows Subsystem for Linux).
You started KubeMQ with Try KubeMQ or Install with Docker.
Turn the Google Pub/Sub connector on
The connector is off by default. This command turns it on and publishes port 8085 on the same container and volume. If you started KubeMQ with the Docker command, skip it and use Without kmq below.
kmq deploy update --installation kubemq --connector gcp --deadline 10mCheck that installation_id is the same as before and data_retained is true. You should see:
{
"installation_id": "…",
"state": "…",
"target": "…",
"data_retained": true,
"next_action": "…"
}Point your client at KubeMQ
export GOOGLE_APPLICATION_CREDENTIALS=/path/to/key.json
export GOOGLE_CLOUD_PROJECT=my-projectexport PUBSUB_EMULATOR_HOST=localhost:8085
export PUBSUB_PROJECT_ID=my-project # any value; the connector ignores the project
unset GOOGLE_APPLICATION_CREDENTIALSSend and receive one message
Install the client library, save the script, and run it. Each run creates and deletes its own topic and subscription.
python3 -m pip install google-cloud-pubsubimport os, time, uuid
from google.cloud import pubsub_v1 # google-cloud-pubsub 2.x
os.environ["PUBSUB_EMULATOR_HOST"] = "localhost:8085"
os.environ["PUBSUB_PROJECT_ID"] = "smoke-test"
suffix = uuid.uuid4().hex[:8] # new names on every run
pub = pubsub_v1.PublisherClient()
sub = pubsub_v1.SubscriberClient()
t = pub.topic_path("smoke-test", f"smoke-topic-{suffix}") # -> gcp.smoke-topic-<suffix>
s = sub.subscription_path("smoke-test", f"smoke-sub-{suffix}") # -> gcp.sub.smoke-sub-<suffix>
pub.create_topic(request={"name": t})
sub.create_subscription(request={"name": s, "topic": t})
pub.publish(t, data=b"smoke-payload", env="ci").result()
print("published")
time.sleep(0.5) # allow fan-out
resp = sub.pull(request={"subscription": s, "max_messages": 1})
assert len(resp.received_messages) == 1, "expected 1 message"
rm = resp.received_messages[0]
assert rm.message.data == b"smoke-payload"
assert rm.message.attributes["env"] == "ci"
sub.acknowledge(request={"subscription": s, "ack_ids": [rm.ack_id]})
print("smoke test PASSED — message received and acked")
sub.delete_subscription(request={"subscription": s})
pub.delete_topic(request={"topic": t})
sub.close()
print("cleaned up")python3 smoke_test.pyYou should see:
published
smoke test PASSED — message received and acked
cleaned upTo turn the connector off again, run the same kmq deploy update command without --connector. Your messages and your evaluation or license stay.
You installed KubeMQ with Install on Kubernetes.
Turn the connector on
Put the connector's key in cluster-values.yaml. If Helm installed your cluster, add it to the file you created on Install on Kubernetes; if kmq installed it, create the file with only this key.
gcp:
enabled: true
sessionAffinity: ClientIPsessionAffinity: ClientIP keeps each client on one server; see Service exposure & session affinity.
If kmq installed your cluster, apply the file with kmq:
kmq deploy update --installation messaging --kube-context YOUR_KUBE_CONTEXT --values cluster-values.yaml --deadline 10mReplace:
YOUR_KUBE_CONTEXT— your cluster's kubectl context;kubectl config get-contextslists them.
Check that data_retained is true. You should see:
{
"installation_id": "…",
"state": "…",
"target": "…",
"data_retained": true,
"next_action": "…"
}If Helm installed your cluster, apply the file with the command in Apply a settings change.
Either way, check that the connector's Service exists:
kubectl --context YOUR_KUBE_CONTEXT -n kubemq get svc messaging-gcpCheck that the row lists 8085/TCP. You should see:
NAME TYPE CLUSTER-IP EXTERNAL-IP PORT(S) AGE
messaging-gcp ClusterIP … <none> 8085/TCP …Reach the connector from your machine
Keep this running in a second terminal:
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-gcp 8085:8085You should see:
Forwarding from 127.0.0.1:8085 -> 8085
Forwarding from [::1]:8085 -> 8085Send and receive one message
Run steps 2 and 3 of the Docker tab unchanged. The port-forward serves the connector at the same localhost:8085 address.
Before you move production traffic
- Read What does not migrate before you commit.
- Plan how traffic moves and how you roll back: run both brokers, move one subscription at a time, and keep Google Cloud Pub/Sub until you have verified the move.
- Install KubeMQ for real: Choose your path.
Compatibility Matrix
| Dimension | Support | Notes |
|---|---|---|
| Drop-in level | endpoint-only | PUBSUB_EMULATOR_HOST=host:8085; no code change |
| Point-to-point queues | N/A | Subscriptions map to KubeMQ Queues; no separate queue API |
| Pub/sub (non-durable) | ✅ | topic → Events Store gcp.{topic}; fan-out per subscription |
| Durable / persistent subscriptions | ✅ | subscription → Queue gcp.sub.{subscription} (broker-durable) |
| Request / reply (RPC) | N/A | Pub/Sub has no RPC primitive |
| Ordering guarantee | ⚠️ node-local | Per ordering key, at-most-one-in-flight; not cluster-wide; lost on process restart |
| Transactions | N/A | Pub/Sub has no transaction concept |
| Dead-letter / redrive | ✅ | Connector-level: re-publishes to the dead-letter topic (new message IDs) when max_delivery_attempts is exceeded — see footnote ¹ |
| Selectors / filtering / wildcards | ✅ | CEL-subset filter on subscription attributes; attributes:K, attributes.K="v", hasPrefix, AND/OR/NOT |
| Auth model | ⚠️ KubeMQ token | When KubeMQ authentication is on, a KubeMQ token in the authorization metadata; no Google IAM or OAuth2; IAM RPCs are permissive stubs |
| TLS / mTLS | ❌ none | Plaintext gRPC only; terminate TLS at a reverse proxy or service mesh |
| Top unsupported | — | No TLS; IAM stubs; BigQuery/GCS export subs; ordering/exactly-once node-local; no default retention |
¹ Dead-lettering is connector-level, not a broker redrive. When a message's receive count exceeds
max_delivery_attempts, the connector republishes it to the configured dead-letter topic through the normal fan-out path. This assigns new message IDs and resets the delivery counter; the broker's own queue-redrive path is not involved.max_delivery_attemptsmust be 5..100; leaving it unset disables the dead-letter policy.
Connection / Endpoint Migration
The only required change is the endpoint env var. The SDK clears Google credentials and uses the insecure emulator path automatically.
Before (real Google Cloud)
# Application authenticates with ADC or a service-account key.
export GOOGLE_APPLICATION_CREDENTIALS=/path/to/key.json
export GOOGLE_CLOUD_PROJECT=my-projectfrom google.cloud import pubsub_v1
publisher = pubsub_v1.PublisherClient()
subscriber = pubsub_v1.SubscriberClient()After (KubeMQ)
# Point every Google Pub/Sub SDK at KubeMQ.
# Clear credentials so the library enters unauthenticated emulator mode.
export PUBSUB_EMULATOR_HOST=kubemq-host:8085
export PUBSUB_PROJECT_ID=my-project # arbitrary; the connector ignores the project segment
unset GOOGLE_APPLICATION_CREDENTIALS
unset GOOGLE_CLOUD_PROJECTNo other client-code change is needed. The connector listens on port 8085 — the same port
Google's emulator uses — so an existing PUBSUB_EMULATOR_HOST that already points at the emulator
only needs its host changed. The connector validates and strips the projects/{p}/ prefix from all
resource paths but otherwise treats resource IDs as global (single-tenant, like the emulator).
Optional — enable AdvertisedEndpoint so the dashboard shows the correct PUBSUB_EMULATOR_HOST
hint:
[Connectors.Gcp]
Enable = true
Port = "8085"
AdvertisedEndpoint = "kubemq.mycompany.svc:8085"Concept & Destination Mapping
| Pub/Sub concept | KubeMQ object | Channel name |
|---|---|---|
Topic projects/{p}/topics/{t} | Events Store log | gcp.{t} |
Subscription projects/{p}/subscriptions/{s} | Queue channel | gcp.sub.{s} |
Message attributes | KubeMQ Tags | passed through as tags |
Message ordering_key | Per-key ordered Queue delivery | embedded in lease metadata |
| Snapshot / Schema | internal registry record | — (no KubeMQ channel) |
Fan-out model. A Publish call writes the message exactly once to the Events Store log
gcp.{t}, then fans out one Queue message per subscription (applying each subscription's filter at
publish time). This means:
- A Pub/Sub topic publish is immediately visible to native KubeMQ consumers on Events Store channel
gcp.{t}. - A subscription's unconsumed backlog is a native KubeMQ Queue on channel
gcp.sub.{s}.
Resource ID rules (same as Google's): 3..255 chars, must start with a letter, charset
[A-Za-z0-9._~%+-], no goog prefix, topic IDs may not start with sub. (reserved broker
namespace).
Message attribute pass-through
A PubsubMessage with attributes and an ordering_key arrives at native consumers with three
reserved tags added:
| Tag | Value |
|---|---|
_pubsub_message_id | Server-assigned message ID |
_pubsub_publish_time | Publish timestamp |
_pubsub_ordering_key | Ordering key (empty string if unset) |
These tags are stripped from attributes when the message is delivered back to a Pub/Sub SDK client.
Canonical Client Example
Client: google-cloud-pubsub 2.x (Python pubsub_v1 gRPC stubs).
import os
from google.cloud import pubsub_v1 # google-cloud-pubsub 2.x
# ── Environment ─────────────────────────────────────────────────────────────
os.environ["PUBSUB_EMULATOR_HOST"] = "localhost:8085"
os.environ["PUBSUB_PROJECT_ID"] = "my-project"
PROJECT = "my-project"
TOPIC = "orders"
SUB = "orders-sub"
# ── Create topic and subscription ────────────────────────────────────────────
publisher = pubsub_v1.PublisherClient() # honours PUBSUB_EMULATOR_HOST
subscriber = pubsub_v1.SubscriberClient()
topic_path = publisher.topic_path(PROJECT, TOPIC) # -> gcp.orders
sub_path = subscriber.subscription_path(PROJECT, SUB) # -> gcp.sub.orders-sub
publisher.create_topic(request={"name": topic_path})
subscriber.create_subscription(
request={"name": sub_path, "topic": topic_path}
)
# ── Publish ───────────────────────────────────────────────────────────────────
future = publisher.publish( # pubsub_v1.PublisherClient.publish
topic_path,
data=b"order-001", # bytes
region="eu-west-1", # arbitrary attributes (passed as Tags)
)
print("published:", future.result()) # blocks until the broker acks
# ── Pull (synchronous) ────────────────────────────────────────────────────────
response = subscriber.pull( # pubsub_v1.SubscriberClient.pull
request={"subscription": sub_path, "max_messages": 10}
)
for msg in response.received_messages:
print("received:", msg.message.data, msg.message.attributes)
subscriber.acknowledge( # pubsub_v1.SubscriberClient.acknowledge
request={"subscription": sub_path, "ack_ids": [msg.ack_id]}
)
# ── StreamingPull (async callback) ───────────────────────────────────────────
def callback(message: pubsub_v1.subscriber.message.Message) -> None:
print("streaming:", message.data, message.attributes)
message.ack() # pubsub_v1.subscriber.message.Message.ack
streaming_future = subscriber.subscribe(sub_path, callback=callback)
try:
streaming_future.result(timeout=10)
except Exception:
streaming_future.cancel()
streaming_future.result()
subscriber.close()RPC is not a Pub/Sub primitive
Pub/Sub has no request/reply mechanism. If you need RPC, use KubeMQ's native Commands/Queries pattern over gRPC or REST rather than modelling it with a reply topic.
Security
Google IAM and OAuth2 are not modelled; the IAM calls (GetIamPolicy, SetIamPolicy,
TestIamPermissions) are permissive stubs. When KubeMQ authentication is on, every call must carry
a KubeMQ token in the authorization metadata (a Bearer prefix is accepted). When it is off, any
client that reaches the port is trusted. The listener has no TLS.
- TLS. For encrypted transport, terminate TLS at a reverse proxy or service mesh in front of port 8085.
- IAM stubs.
GetIamPolicyreturns an emptyPolicy{Version: 3}. No permissions are checked. - DoS guards remain active.
MaxRecvMsgSize, the per-subscription in-flight cap (MaxInflightPerSubscription),MaxConcurrentPolls,MaxSeekReplay, and push delivery backoff are enforced regardless of auth mode.
Security: do not expose port 8085 to untrusted networks
Without KubeMQ authentication, any client that can reach the port can create, publish to, and delete any topic or subscription. Keep the listener inside a trusted network boundary, and terminate TLS at a proxy or mesh if you need encryption.
What Does NOT Migrate / Deviations
Features not supported (hard rejections)
| Feature | Behavior |
|---|---|
| Google IAM / TLS | No Google IAM or OAuth2; IAM is a permissive stub. No TLS on the listener. KubeMQ token authentication applies when it is on (see Security). |
| BigQuery / Cloud Storage / Bigtable export subscriptions | Rejected with INVALID_ARGUMENT — no KubeMQ analog. |
| Ingestion sources (Kinesis, Cloud Storage, Azure Event Hubs) | Rejected with INVALID_ARGUMENT. |
KMS key names (kms_key_name) | Accepted and silently ignored. |
| gRPC REST/JSON (grpc-gateway) | gRPC only; no https://pubsub.googleapis.com/v1/… REST surface. |
Behavioral deviations
| Area | Google Pub/Sub behavior | KubeMQ behavior |
|---|---|---|
| Ordering / exactly-once scope | Cluster-wide | Node-local. An ack_id is only valid on the node that minted it; ordering sequences and exactly-once guarantees are lost on process restart or if the client reconnects to a different cluster node. Pin a subscription's StreamingPull traffic to one node, or accept at-least-once across nodes. |
| Retention default | 7-day default when unset | No connector-level default. Retention is unset unless the client supplies message_retention_duration. Supplied values are bounded to 10 min..31 days and then clamped down to the server's Store.MaxRetention ceiling (when that ceiling is non-zero). There is no 24-hour default. |
| Dead-letter message IDs | Original message IDs preserved | Republishes with new message IDs; resets the counter; max_delivery_attempts 5..100 — see footnote ¹. |
| Push OIDC token | Signed by Google | Signed by the emulator (not Google-verifiable). |
| CEL filter | Full CEL | Attributes-only subset: attributes:K, attributes.K="v", hasPrefix(attributes.K, "p"), AND/OR/NOT/-. data and metadata fields are not filterable. Malformed expressions → INVALID_ARGUMENT. |
| Exactly-once unary ack status | FAILED_PRECONDITION + ErrorInfo(PERMANENT_FAILURE_INVALID_ACK_ID) | Same (matches the real SDK contract). |
DeleteTopic semantics | Topic and its subscriptions are deleted | Tombstone only — the Events Store log is retained so existing subscriptions survive; re-creating the topic reuses the log. |
| Resource ID namespace | Scoped by project | Global (single-tenant). Project is parsed and validated but ignored. |
Seek replay cap | Unlimited | Capped at MaxSeekReplay (default 1 000 000 messages). Hitting the cap stops with a WARN; no silent loss. |
Seek to timestamp before retention window | Error | Clamped to the earliest retained message — not an error. |
Verification Smoke Test
For the basic round trip, see the quick start; this section adds deeper checks.
Cross-protocol fan-out check
The headline behavior of the connector is that a Pub/Sub publish is also a native KubeMQ
message. Create a topic smoke-topic and a subscription smoke-sub with the same client calls
the quick start uses, but with fixed names and no cleanup, so a native KubeMQ client can read them
back:
t = pub.topic_path("smoke-test", "smoke-topic")
s = sub.subscription_path("smoke-test", "smoke-sub")
pub.create_topic(request={"name": t})
sub.create_subscription(request={"name": s, "topic": t})A publish to this topic lands on the Events Store log gcp.smoke-topic, and the backlog of this
subscription is the native Queue gcp.sub.smoke-sub. A native KubeMQ client confirms the fan-out
with no Pub/Sub SDK involved:
# Events Store log carries the topic publish:
#
# SubscribeToEventsStore(channel="gcp.smoke-topic", startAt="new")
#
# Subscription backlog is a native Queue channel:
#
# ReceiveQueueMessages(channel="gcp.sub.smoke-sub", maxMessages=10)For a deterministic read, subscribe to the Events Store log with start policy startAt = "new"
before publishing. See Channel Mapping
for the full gcp.{topic} / gcp.sub.{subscription} scheme.
See Also
Migrate from another broker
Pick the right KubeMQ wire-protocol connector for your broker and migrate onto it.
Getting Started
Point the SDK at port 8085 via PUBSUB_EMULATOR_HOST and run your first round-trip.
Architecture
The gRPC emulator listener, the RPC surface, and how topics and subscriptions map to KubeMQ primitives.
Channel Mapping
The gcp.{topic} / gcp.sub.{subscription} mapping — no rename needed when migrating.
Capabilities
Supported RPCs, out-of-scope operations, and the documented deviations from Google Pub/Sub.
Limits & Rules
Resource-ID rules, retention clamping, max_delivery_attempts bounds, and the seek replay cap.
Error Codes
The INVALID_ARGUMENT / FAILED_PRECONDITION mapping and the ErrorInfo reasons.
Configuration reference
All [Connectors.Gcp] fields, environment variables, and validation rules.
Events Store
The durable pub/sub pattern backing Pub/Sub topics (gcp.{topic}).
Queues
The point-to-point pattern backing Pub/Sub subscriptions (gcp.sub.{subscription}).
Migrating from AWS SQS/SNS
The comparable emulator-endpoint migration for AWS SQS and SNS.
Was this page helpful?