The Go client takes a context.Context on every call and separates building from executing: a
chain is inert until Execute(ctx), Pop(ctx) or Get(ctx). It needs Go 1.24, and its only
direct dependencies are github.com/google/uuid and github.com/jackc/pgx/v5.
go get github.com/smartpricing/queen/clients/client-goimport queen "github.com/smartpricing/queen/clients/client-go"
client, err := queen.New(queen.ClientConfig{
URLs: []string{"http://broker-a:6632", "http://broker-b:6632"},
BearerToken: os.Getenv("QUEEN_TOKEN"),
EnableFailover: true,
})
if err != nil {
return err
}
defer client.Close(context.Background())New also accepts a plain URL string or a slice of them, and returns an error instead of
panicking. Zero-valued fields fall back to ClientDefaults, with one exception: EnableFailover
is false on a bare struct, so set it yourself whenever you pass more than one URL.
Push
res, err := client.Queue("orders").
Partition("customer-42").
Push(map[string]any{"orderId": 9137, "amount": 99.5}).
Execute(ctx)
if err != nil {
return err
}
// res[0].Status == "queued"Push takes one payload, a []interface{} or a []map[string]interface{}, and Execute returns
one PushResponse per item in request order, each carrying the broker’s status. Chain
.TransactionID(id) before Execute and a retry inside the dedup window writes nothing a second
time.
Consume
err = client.Queue("orders").
Group("billing").
SubscriptionMode("all").
Limit(1).
Each().
Consume(ctx, func(ctx context.Context, msg *queen.Message) error {
fmt.Println(msg.Data)
return nil
}).
Execute(ctx)
if err != nil {
return err
}Limit(1) is what ends that loop: without it Execute long-polls for the next message. A worker
is the same chain with Concurrency(4) and no limit, and ConsumeBatch hands the handler the
whole claimed batch instead of one message at a time.
Acknowledge
AutoAck is on by default and is client-side: the worker acks completed when the handler
returns nil, failed when it returns an error. An ack is an offset commit, so a nack clamps the
cursor and every message after the failed one in the batch comes back, which is why .Each()
abandons the rest of it.
Turn AutoAck off and settle the batch yourself:
responses, err := client.Ack(ctx, messages, true, queen.AckOptions{ConsumerGroup: "billing"})Ack accepts *Message, Message, []*Message or []Message, and success maps to
completed or failed. A rejected ack still arrives as HTTP 200: check responses[i].Success, a
nil error is not enough.
What differs here
- A bare
Pop(ctx)does not long-poll: call.Wait(true). RetryAttemptscounts retries after the first attempt, where the other SDKs count total tries.TransactionIDapplies to the first pushed item only; the rest get fresh UUIDv7s.Message.Queue,RetryCountandErrorMessagestay zero after a pop, and are filled on a DLQ read.- No signal handlers are installed: pair
defer client.Close(ctx)with your ownsignal.NotifyContext.
Every client option, the full consume table, transactions, buffering, the dead-letter queue, the admin surface and the streaming SDK are in the Go reference.
Tutorials
Five programs that build on each other, each one asserting its own outcome against a live broker: one message end to end, a multi-queue flow, transactional handoff, replay, and a streaming aggregation.
Start with Hello world.