# Stream API (Upstream/Downstream) (/learn/queues/tutorials/stream-api)



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

A streaming producer that sends messages continuously through an upstream connection, and a streaming consumer that receives messages through a downstream connection with range-based acknowledgment.

<Mermaid
  chart="graph LR
  P[&#x22;Producer&#x22;]
  UP{{&#x22;Upstream stream<br/>(send)&#x22;}}
  Q[[&#x22;Queue<br/>orders.stream&#x22;]]
  DOWN{{&#x22;Downstream stream<br/>(receive)&#x22;}}
  C[&#x22;Consumer&#x22;]

  P -- &#x22;send batch&#x22; --> UP
  UP -- &#x22;enqueue&#x22; --> Q
  Q -- &#x22;deliver batch&#x22; --> DOWN
  DOWN -- &#x22;process&#x22; --> C
  C -. &#x22;ack range&#x22; .-> DOWN

  class UP,DOWN stream
  class Q queue
  class P,C client"
/>

*One long-lived upstream stream feeds the queue; a downstream stream delivers batches that the consumer settles with a dotted range acknowledgment.*

## Prerequisites [#prerequisites]

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

## Upstream: Stream Send [#upstream-stream-send]

Open a persistent bidirectional stream for high-throughput sending. Each message gets an individual response confirming receipt.

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="upstream.go"
    stream, err := client.UpstreamQueue(ctx)
    if err != nil {
        log.Fatal(err)
    }

    for i := 0; i < 100; i++ {
        body := fmt.Sprintf(`{"orderId":"ORD-%04d","item":"widget"}`, i)
        result, err := stream.Send(ctx, kubemq.NewQueueMessage().
            SetChannel("orders.stream").
            SetBody([]byte(body)),
        )
        if err != nil {
            log.Printf("Send error: %v", err)
            continue
        }
        fmt.Printf("Streamed: id=%s\n", result.MessageID)
    }
    stream.Close()
    ```
  </Tab>

  <Tab value="Python">
    ```python title="upstream.py"
    stream = client.upstream_queue()

    for i in range(100):
        result = stream.send(
            QueueMessage(
                channel="orders.stream",
                body=f'{{"orderId":"ORD-{i:04d}","item":"widget"}}'.encode(),
            )
        )
        print(f"Streamed: id={result.id}")

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

  <Tab value="Node.js">
    ```typescript title="upstream.ts"
    const upstream = client.createQueueUpstream();

    for (let i = 0; i < 100; i++) {
      const result = await upstream.send([
        createQueueMessage({
          channel: 'orders.stream',
          body: JSON.stringify({ orderId: `ORD-${String(i).padStart(4, '0')}`, item: 'widget' }),
        }),
      ]);
      console.log(`Streamed: id=${result.results[0].messageId}`);
    }

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

  <Tab value="Java">
    ```java title="Upstream.java"
    QueueUpstream stream = client.upstreamQueue();

    for (int i = 0; i < 100; i++) {
        SendQueueMessageResult result = stream.send(QueueMessage.builder()
            .channel("orders.stream")
            .body(String.format("{\"orderId\":\"ORD-%04d\",\"item\":\"widget\"}", i).getBytes())
            .build());
        System.out.printf("Streamed: id=%s%n", result.getMessageId());
    }

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

  <Tab value="C#">
    ```csharp title="Upstream.cs"
    var stream = await client.UpstreamQueueAsync();

    for (int i = 0; i < 100; i++)
    {
        var result = await stream.SendAsync(new QueueMessage
        {
            Channel = "orders.stream",
            Body = Encoding.UTF8.GetBytes($"{{\"orderId\":\"ORD-{i:D4}\",\"item\":\"widget\"}}")
        });
        Console.WriteLine($"Streamed: id={result.MessageId}");
    }

    await stream.CloseAsync();
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="Upstream.kt"
    val stream = client.upstreamQueue()

    for (i in 0 until 100) {
        val result = stream.send(QueueMessage(
            channel = "orders.stream",
            body = """{"orderId":"ORD-${"%04d".format(i)}","item":"widget"}""".toByteArray()
        ))
        println("Streamed: id=${result.messageId}")
    }

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

  <Tab value="C++">
    ```cpp title="upstream.cpp"
    auto stream = client.upstreamQueue();

    for (int i = 0; i < 100; i++) {
        kubemq::QueueMessage msg;
        msg.channel = "orders.stream";
        msg.body = "{\"orderId\":\"ORD-" + std::to_string(i) + "\",\"item\":\"widget\"}";
        auto result = stream.send(msg);
        std::cout << "Streamed: id=" << result.messageId << std::endl;
    }

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

  <Tab value="Rust">
    ```rust title="upstream.rs"
    let mut upstream = client.queue_upstream().await?;

    // Send a batch of messages over the persistent upstream stream.
    let messages: Vec<QueueMessage> = (0..100)
        .map(|i| {
            QueueMessageBuilder::new()
                .channel("orders.stream")
                .body(format!(r#"{{"orderId":"ORD-{:04}","item":"widget"}}"#, i).into_bytes())
                .build()
        })
        .collect();

    upstream.send("orders-batch-001", messages).await?;

    if let Some(result) = upstream.results().recv().await {
        println!(
            "Streamed: ref_id={}, is_error={}, items={}",
            result.ref_request_id,
            result.is_error,
            result.results.len()
        );
    }

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

  <Tab value="Ruby">
    ```ruby title="upstream.rb"
    sender = client.create_upstream_sender

    100.times do |i|
      msg = KubeMQ::Queues::QueueMessage.new(
        channel: "orders.stream",
        body: %({"orderId":"ORD-#{format('%04d', i)}","item":"widget"})
      )
      results = sender.publish(msg)
      results.each { |r| puts "Streamed: id=#{r.id}, error?=#{r.error?}" }
    end

    sender.close
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="upstream.exs"
    {:ok, handle} = KubeMQ.Client.queue_upstream(client)

    # Send a batch of messages over the persistent upstream stream.
    messages =
      for i <- 0..99 do
        body = ~s({"orderId":"ORD-#{String.pad_leading(Integer.to_string(i), 4, "0")}","item":"widget"})
        KubeMQ.QueueMessage.new(channel: "orders.stream", body: body)
      end

    case KubeMQ.QueueUpstreamHandle.send(handle, messages) do
      {:ok, results} ->
        Enum.each(results, fn r ->
          IO.puts("Streamed: id=#{r.message_id}, error: #{r.is_error}")
        end)

      {:error, err} ->
        IO.puts("Stream send failed: #{err.message}")
    end

    KubeMQ.QueueUpstreamHandle.close(handle)
    ```
  </Tab>
</Tabs>

## Downstream: Stream Receive [#downstream-stream-receive]

Open a persistent downstream stream for continuous message consumption.

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="downstream.go"
    stream, err := client.DownstreamQueue(ctx, &kubemq.DownstreamRequest{
        Channel:            "orders.stream",
        MaxItems:           10,
        WaitTimeoutSeconds: 5,
        AutoAck:            false,
    })
    if err != nil {
        log.Fatal(err)
    }

    for resp := range stream.ResponseCh {
        for _, m := range resp.Messages {
            fmt.Printf("Received: %s\n", string(m.Message.Body))
        }
        resp.AckAll()
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python title="downstream.py"
    def on_message(response):
        for msg in response.messages:
            print(f"Received: {msg.body.decode('utf-8')}")
            msg.ack()

    def on_error(err):
        print(f"Stream error: {err}")

    stream = client.downstream_queue(
        channel="orders.stream",
        max_messages=10,
        wait_timeout_in_seconds=5,
        on_message_callback=on_message,
        on_error_callback=on_error,
    )
    ```
  </Tab>

  <Tab value="Node.js">
    ```typescript title="downstream.ts"
    const stream = client.streamQueueMessages({
      channel: 'orders.stream',
      maxMessages: 10,
      waitTimeoutSeconds: 5,
      autoAck: false,
    });

    stream.onMessages((messages) => {
      for (const msg of messages) {
        console.log('Received:', new TextDecoder().decode(msg.body));
      }
      // Acknowledge the whole batch at once.
      stream.ackAll();
    });

    stream.onError((err) => console.error('Stream error:', err.message));
    ```
  </Tab>

  <Tab value="Java">
    ```java title="Downstream.java"
    client.downstreamQueue(DownstreamRequest.builder()
        .channel("orders.stream")
        .maxMessages(10)
        .waitTimeoutSeconds(5)
        .onMessage(response -> {
            for (QueueMessageReceived msg : response.getMessages()) {
                System.out.println("Received: " + new String(msg.getBody()));
                msg.ack();
            }
        })
        .onError(err -> System.err.println("Stream error: " + err.getMessage()))
        .build());
    ```
  </Tab>

  <Tab value="C#">
    ```csharp title="Downstream.cs"
    await foreach (var response in client.DownstreamQueueAsync(new DownstreamRequest
    {
        Channel = "orders.stream",
        MaxMessages = 10,
        WaitTimeoutSeconds = 5,
    }))
    {
        foreach (var msg in response.Messages)
        {
            Console.WriteLine($"Received: {Encoding.UTF8.GetString(msg.Body.Span)}");
            await msg.AckAsync();
        }
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="Downstream.kt"
    client.downstreamQueue(
        channel = "orders.stream",
        maxMessages = 10,
        waitTimeoutSeconds = 5,
        onMessage = { response ->
            for (msg in response.messages) {
                println("Received: ${String(msg.body)}")
                msg.ack()
            }
        },
        onError = { err -> System.err.println("Stream error: ${err.message}") }
    )
    ```
  </Tab>

  <Tab value="C++">
    ```cpp title="downstream.cpp"
    client.downstreamQueue("orders.stream", 10, 5,
        [](const auto& response) {
            for (const auto& msg : response.messages) {
                std::cout << "Received: " << msg.body << std::endl;
                msg.ack();
            }
        },
        [](const std::string& err) {
            std::cerr << "Stream error: " << err << std::endl;
        }
    );
    ```
  </Tab>

  <Tab value="Rust">
    ```rust title="downstream.rs"
    let mut receiver = client.new_queue_downstream_receiver().await?;

    // Poll a batch over the persistent downstream stream.
    let poll = PollRequest {
        channel: "orders.stream".to_string(),
        max_items: 10,
        wait_timeout_seconds: 5,
        auto_ack: false,
    };

    let response = receiver.poll(poll).await?;
    for msg in &response.messages {
        println!("Received: {}", String::from_utf8_lossy(&msg.body));
    }

    // Acknowledge the whole batch (range ack) in one call.
    response.ack_all().await?;

    receiver.close().await?;
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby title="downstream.rb"
    receiver = client.create_downstream_receiver
    request = KubeMQ::Queues::QueuePollRequest.new(
      channel: "orders.stream",
      max_items: 10,
      wait_timeout: 5
    )
    response = receiver.poll(request)

    if response.error?
      puts "Stream error: #{response.error}"
    else
      response.messages.each do |msg|
        puts "Received: #{msg.body}"
        msg.ack
      end
    end

    receiver.close
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="downstream.exs"
    # Poll a batch over the downstream stream.
    case KubeMQ.Client.poll_queue(client,
           channel: "orders.stream",
           max_items: 10,
           wait_timeout: 5_000
         ) do
      {:ok, poll} ->
        Enum.each(poll.messages, fn msg ->
          IO.puts("Received: #{msg.body}")
        end)

        # Acknowledge the whole batch (range ack) in one call.
        {:ok, _} = KubeMQ.PollResponse.ack_all(poll)

      {:error, err} ->
        IO.puts("Stream error: #{err.message}")
    end
    ```
  </Tab>
</Tabs>

## Stream vs Polling Comparison [#stream-vs-polling-comparison]

| Feature        | Stream API                        | Polling (PollQueue)            |
| -------------- | --------------------------------- | ------------------------------ |
| Connection     | Persistent bidirectional          | Request/response per poll      |
| Latency        | Lower (always connected)          | Higher (new request each time) |
| Throughput     | Higher                            | Lower                          |
| Resource usage | Holds connection open             | Releases between polls         |
| Best for       | High-volume continuous processing | Periodic batch processing      |

## Next Steps [#next-steps]

<Cards>
  <Card title="Batch Operations" href="/learn/queues/tutorials/batch-operations" description="Batch send and receive for high throughput." />

  <Card title="Peek Messages" href="/learn/queues/tutorials/peek-messages" description="Inspect queue contents without consuming." />
</Cards>
