# Send & Receive Messages (/learn/queues/tutorials/send-receive)



## What You Will Build [#what-you-will-build]

A producer that sends order processing tasks with metadata and tags, and a consumer that receives, processes, and acknowledges each message. The example includes error handling and message introspection.

<Mermaid
  chart="sequenceDiagram
    participant P as Producer
    participant Q as KubeMQ Queue
    participant C as Consumer
    P->>Q: Send message (body + metadata + tags)
    Q-->>P: MessageID + SentAt
    C->>Q: Poll (maxItems, waitTimeout)
    Q-->>C: Deliver message
    Note over C: Process message
    C->>Q: Ack
    Note over Q: Message removed"
/>

*The send-receive-acknowledge cycle: the producer sends, the consumer polls and processes, and the ack removes the message from the queue.*

## Prerequisites [#prerequisites]

* KubeMQ server running on `localhost:50000`
* SDK installed ([Getting Started](../getting-started))

## Steps [#steps]

<Steps>
  <Step>
    ### Set Up the Client [#set-up-the-client]

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go title="main.go"
        package main

        import (
            "context"
            "encoding/json"
            "fmt"
            "log"
            "time"

            "github.com/kubemq-io/kubemq-go/v2"
        )

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

            client, err := kubemq.NewClient(ctx,
                kubemq.WithAddress("localhost", 50000),
                kubemq.WithClientId("order-processor"),
            )
            if err != nil {
                log.Fatal(err)
            }
            defer client.Close()
        ```
      </Tab>

      <Tab value="Python">
        ```python title="main.py"
        from kubemq.queues import Client as QueuesClient
        from kubemq import QueueMessage

        client = QueuesClient(
            address="localhost:50000",
            client_id="order-processor",
        )
        ```
      </Tab>

      <Tab value="Node.js">
        ```typescript title="main.ts"
        import { KubeMQClient, createQueueMessage } from 'kubemq-js';

        const client = await KubeMQClient.create({
          address: 'localhost:50000',
          clientId: 'order-processor',
        });
        ```
      </Tab>

      <Tab value="Java">
        ```java title="Main.java"
        QueuesClient client = QueuesClient.builder()
            .address("localhost:50000")
            .clientId("order-processor")
            .build();
        ```
      </Tab>

      <Tab value="C#">
        ```csharp title="Program.cs"
        using KubeMQ.Sdk.Client;

        await using var client = new KubeMQClient(new KubeMQClientOptions());
        await client.ConnectAsync();
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin title="Main.kt"
        val client = QueuesClient("localhost:50000")
        ```
      </Tab>

      <Tab value="C++">
        ```cpp title="main.cpp"
        #include <kubemq/client.h>
        auto client = kubemq::QueuesClient("localhost:50000");
        ```
      </Tab>

      <Tab value="Rust">
        ```rust title="main.rs"
        use kubemq::prelude::*;
        use kubemq::QueueMessageBuilder;
        use std::collections::HashMap;

        #[tokio::main]
        async fn main() -> kubemq::Result<()> {
            let client = KubemqClient::builder()
                .host("localhost")
                .port(50000)
                .client_id("order-processor")
                .build()
                .await?;
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby title="main.rb"
        require 'kubemq'

        client = KubeMQ::QueuesClient.new(
          address: 'localhost:50000',
          client_id: 'order-processor',
        )
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir title="main.exs"
        {:ok, client} = KubeMQ.Client.start_link(
          address: "localhost:50000",
          client_id: "order-processor"
        )
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Send a Message with Metadata and Tags [#send-a-message-with-metadata-and-tags]

    Create a queue message with a JSON body, metadata string, and tags for downstream routing.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go
            order := map[string]interface{}{
                "orderId": "ORD-5001",
                "items":   3,
                "total":   149.97,
            }
            body, _ := json.Marshal(order)

            msg := kubemq.NewQueueMessage().
                SetChannel("orders").
                SetBody(body).
                SetMetadata("order.created").
                SetTags(map[string]string{
                    "region":   "us-east",
                    "priority": "high",
                })

            result, err := client.SendQueueMessage(ctx, msg)
            if err != nil {
                log.Fatal(err)
            }
            if result.IsError {
                log.Fatalf("Send failed: %s", result.Error)
            }
            fmt.Printf("Sent: id=%s, sentAt=%d\n", result.MessageID, result.SentAt)
        ```
      </Tab>

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

        order = {"orderId": "ORD-5001", "items": 3, "total": 149.97}

        result = client.send_queue_message(
            QueueMessage(
                channel="orders",
                body=json.dumps(order).encode(),
                metadata="order.created",
                tags={"region": "us-east", "priority": "high"},
            )
        )
        print(f"Sent: id={result.id}, sentAt={result.sent_at}")
        ```
      </Tab>

      <Tab value="Node.js">
        ```typescript
        const result = await client.sendQueueMessage(
          createQueueMessage({
            channel: 'orders',
            body: JSON.stringify({ orderId: 'ORD-5001', items: 3, total: 149.97 }),
            metadata: 'order.created',
            tags: { region: 'us-east', priority: 'high' },
          }),
        );
        console.log(`Sent: id=${result.messageId}, sentAt=${result.sentAt}`);
        ```
      </Tab>

      <Tab value="Java">
        ```java
        QueueMessage msg = QueueMessage.builder()
            .channel("orders")
            .body("{\"orderId\":\"ORD-5001\",\"items\":3,\"total\":149.97}".getBytes())
            .metadata("order.created")
            .tags(Map.of("region", "us-east", "priority", "high"))
            .build();

        SendQueueMessageResult result = client.sendQueueMessage(msg);
        System.out.printf("Sent: id=%s, sentAt=%d%n", result.getMessageId(), result.getSentAt());
        ```
      </Tab>

      <Tab value="C#">
        ```csharp
        var result = await client.SendQueueMessageAsync(new QueueMessage
        {
            Channel = "orders",
            Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-5001\",\"items\":3,\"total\":149.97}"),
            Metadata = "order.created",
            Tags = new Dictionary<string, string>
            {
                ["region"] = "us-east",
                ["priority"] = "high"
            }
        });
        Console.WriteLine($"Sent: id={result.MessageId}, sentAt={result.SentAt}");
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin
        val result = client.sendQueueMessage(QueueMessage(
            channel = "orders",
            body = """{"orderId":"ORD-5001","items":3,"total":149.97}""".toByteArray(),
            metadata = "order.created",
            tags = mapOf("region" to "us-east", "priority" to "high")
        ))
        println("Sent: id=${result.messageId}, sentAt=${result.sentAt}")
        ```
      </Tab>

      <Tab value="C++">
        ```cpp
        kubemq::QueueMessage msg;
        msg.channel = "orders";
        msg.body = R"({"orderId":"ORD-5001","items":3,"total":149.97})";
        msg.metadata = "order.created";
        msg.tags = {{"region", "us-east"}, {"priority", "high"}};

        auto result = client.sendQueueMessage(msg);
        std::cout << "Sent: id=" << result.messageId << std::endl;
        ```
      </Tab>

      <Tab value="Rust">
        ```rust
            let order = r#"{"orderId":"ORD-5001","items":3,"total":149.97}"#;
            let mut tags = HashMap::new();
            tags.insert("region".to_string(), "us-east".to_string());
            tags.insert("priority".to_string(), "high".to_string());

            let msg = QueueMessageBuilder::new()
                .channel("orders")
                .body(order.as_bytes().to_vec())
                .metadata("order.created")
                .tags(tags)
                .build();

            let result = client.send_queue_message(msg).await?;
            println!("Sent: id={}, sent_at={}", result.message_id, result.sent_at);
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby
        msg = KubeMQ::Queues::QueueMessage.new(
          channel: 'orders',
          body: '{"orderId":"ORD-5001","items":3,"total":149.97}',
          metadata: 'order.created',
          tags: { 'region' => 'us-east', 'priority' => 'high' }
        )

        result = client.send_queue_message(msg)
        puts "Sent: id=#{result.id}, error?=#{result.error?}"
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir
        msg = KubeMQ.QueueMessage.new(
          channel: "orders",
          body: ~s({"orderId":"ORD-5001","items":3,"total":149.97}),
          metadata: "order.created",
          tags: %{"region" => "us-east", "priority" => "high"}
        )

        {:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
        IO.puts("Sent: id=#{result.message_id}")
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Receive Messages [#receive-messages]

    Poll the queue for available messages. You control how many messages to fetch and how long to wait.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go
            resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
                Channel:            "orders",
                MaxItems:           10,
                WaitTimeoutSeconds: 5,
                AutoAck:            false,
            })
            if err != nil {
                log.Fatal(err)
            }
            fmt.Printf("Received %d messages\n", len(resp.Messages))
        ```
      </Tab>

      <Tab value="Python">
        ```python
        response = client.receive_queue_messages(
            channel="orders",
            max_messages=10,
            wait_timeout_in_seconds=5,
        )
        print(f"Received {len(response.messages)} messages")
        ```
      </Tab>

      <Tab value="Node.js">
        ```typescript
        const messages = await client.receiveQueueMessages({
          channel: 'orders',
          maxMessages: 10,
          waitTimeoutSeconds: 5,
        });
        console.log(`Received ${messages.length} messages`);
        ```
      </Tab>

      <Tab value="Java">
        ```java
        ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
            ReceiveQueueMessagesRequest.builder()
                .channel("orders")
                .maxMessages(10)
                .waitTimeoutSeconds(5)
                .build());

        System.out.printf("Received %d messages%n", response.getMessages().size());
        ```
      </Tab>

      <Tab value="C#">
        ```csharp
        var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
        {
            Channel = "orders",
            MaxMessages = 10,
            WaitTimeoutSeconds = 5,
        });
        Console.WriteLine($"Received {response.Messages.Count} messages");
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin
        val response = client.receiveQueueMessages(
            channel = "orders",
            maxMessages = 10,
            waitTimeoutSeconds = 5
        )
        println("Received ${response.messages.size} messages")
        ```
      </Tab>

      <Tab value="C++">
        ```cpp
        auto response = client.receiveQueueMessages("orders", 10, 5);
        std::cout << "Received " << response.messages.size() << " messages" << std::endl;
        ```
      </Tab>

      <Tab value="Rust">
        ```rust
            // receive_queue_messages(channel, max_items, wait_seconds, auto_ack)
            let messages = client
                .receive_queue_messages("orders", 10, 5, false)
                .await?;
            println!("Received {} messages", messages.len());
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby
        messages = client.receive_queue_messages(
          channel: 'orders',
          max_messages: 10,
          wait_timeout_seconds: 5
        )
        puts "Received #{messages.size} messages"
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir
        {:ok, response} =
          KubeMQ.Client.receive_queue_messages(client, "orders",
            max_messages: 10,
            wait_timeout: 5_000
          )

        IO.puts("Received #{response.messages_received} messages")
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Process and Acknowledge [#process-and-acknowledge]

    Inspect each message's body, metadata, and tags, then acknowledge to remove it from the queue.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go
            for _, dm := range resp.Messages {
                fmt.Printf("  ID:       %s\n", dm.Message.MessageID)
                fmt.Printf("  Body:     %s\n", string(dm.Message.Body))
                fmt.Printf("  Metadata: %s\n", dm.Message.Metadata)
                fmt.Printf("  Tags:     %v\n", dm.Message.Tags)
            }
            if err := resp.AckAll(); err != nil {
                log.Fatal(err)
            }
            fmt.Println("All messages acknowledged")
        }
        ```
      </Tab>

      <Tab value="Python">
        ```python
        for msg in response.messages:
            print(f"  ID:       {msg.id}")
            print(f"  Body:     {msg.body.decode('utf-8')}")
            print(f"  Metadata: {msg.metadata}")
            print(f"  Tags:     {msg.tags}")
            msg.ack()
            print("  Acknowledged")

        client.close()
        ```
      </Tab>

      <Tab value="Node.js">
        ```typescript
        for (const msg of messages) {
          console.log('  ID:      ', msg.messageId);
          console.log('  Body:    ', new TextDecoder().decode(msg.body));
          console.log('  Metadata:', msg.metadata);
          console.log('  Tags:    ', msg.tags);
          await msg.ack();
          console.log('  Acknowledged');
        }

        await client.close();
        ```
      </Tab>

      <Tab value="Java">
        ```java
        for (QueueMessageReceived msg : response.getMessages()) {
            System.out.printf("  ID:       %s%n", msg.getMessageId());
            System.out.printf("  Body:     %s%n", new String(msg.getBody()));
            System.out.printf("  Metadata: %s%n", msg.getMetadata());
            System.out.printf("  Tags:     %s%n", msg.getTags());
            msg.ack();
            System.out.println("  Acknowledged");
        }

        client.close();
        ```
      </Tab>

      <Tab value="C#">
        ```csharp
        foreach (var msg in response.Messages)
        {
            Console.WriteLine($"  ID:       {msg.MessageId}");
            Console.WriteLine($"  Body:     {Encoding.UTF8.GetString(msg.Body.Span)}");
            Console.WriteLine($"  Metadata: {msg.Metadata}");
            Console.WriteLine($"  Tags:     {string.Join(", ", msg.Tags)}");
            await msg.AckAsync();
            Console.WriteLine("  Acknowledged");
        }
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin
        for (msg in response.messages) {
            println("  ID:       ${msg.messageId}")
            println("  Body:     ${String(msg.body)}")
            println("  Metadata: ${msg.metadata}")
            println("  Tags:     ${msg.tags}")
            msg.ack()
            println("  Acknowledged")
        }

        client.close()
        ```
      </Tab>

      <Tab value="C++">
        ```cpp
        for (const auto& msg : response.messages) {
            std::cout << "  ID:       " << msg.messageId << std::endl;
            std::cout << "  Body:     " << msg.body << std::endl;
            std::cout << "  Metadata: " << msg.metadata << std::endl;
            msg.ack();
            std::cout << "  Acknowledged" << std::endl;
        }
        ```
      </Tab>

      <Tab value="Rust">
        ```rust
            // The unary receive returns plain messages; settle them with ack_all.
            // For per-message ack/nack, use the queue stream API.
            for m in &messages {
                println!("  ID:       {}", m.id);
                println!("  Body:     {}", String::from_utf8_lossy(&m.body));
                println!("  Metadata: {}", m.metadata);
                println!("  Tags:     {:?}", m.tags);
            }

            let ack = AckAllQueueMessagesRequest {
                request_id: String::new(),
                client_id: String::new(),
                channel: "orders".to_string(),
                wait_time_seconds: 5,
            };
            client.ack_all_queue_messages(&ack).await?;
            println!("All messages acknowledged");

            client.close().await?;
            Ok(())
        }
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby
        messages.each do |m|
          puts "  ID:       #{m.id}"
          puts "  Body:     #{m.body}"
          puts "  Metadata: #{m.metadata}"
          puts "  Tags:     #{m.tags}"
        end

        # Settle the polled messages. For per-message ack/nack, use the stream receiver.
        affected = client.ack_all_queue_messages(channel: 'orders', wait_timeout_seconds: 5)
        puts "Acknowledged #{affected} messages"

        client.close
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir
        Enum.each(response.messages, fn m ->
          IO.puts("  ID:       #{m.id}")
          IO.puts("  Body:     #{m.body}")
          IO.puts("  Metadata: #{m.metadata}")
          IO.puts("  Tags:     #{inspect(m.tags)}")
        end)

        # Settle the received messages. For per-message ack/nack, use poll_queue (stream API).
        {:ok, ack} = KubeMQ.Client.ack_all_queue_messages(client, "orders", wait_timeout: 5_000)
        IO.puts("Acknowledged #{ack.affected_messages} messages")

        KubeMQ.Client.close(client)
        ```
      </Tab>
    </Tabs>
  </Step>
</Steps>

## Next Steps [#next-steps]

<Cards>
  <Card title="Ack, Nack & Requeue" href="/learn/queues/tutorials/ack-nack-requeue" description="Master the three message settlement options." />

  <Card title="Dead Letter Queue" href="/learn/queues/tutorials/dead-letter-queue" description="Route failed messages to a DLQ after max retries." />

  <Card title="Batch Operations" href="/learn/queues/tutorials/batch-operations" description="Send and receive multiple messages in one request." />
</Cards>
