Skip to content

Send Message ​

The Send Message block publishes messages from a route or workflow to a message queue: a Kafka topic, a NATS JetStream subject, an Amazon SQS queue, a RabbitMQ exchange or queue, or a Redis stream. It waits until the broker (the queue server) confirms each message, then tells you what was sent.

When to use it ​

  • Hand work to another service: "order created", "send this email", "resize this image".
  • Fan an event out to everything that listens on a topic.
  • Feed a trigger that starts another workflow from the same queue.
  • Don't use it to start one of your own workflows directly. Trigger Workflow does that with no queue to set up.

Editions

Sending to RabbitMQ or a Redis stream works in every edition. Kafka, NATS and SQS need an enterprise license, the same as their triggers. Saving a block that points at one of them without a license is refused with the same error. A block saved while licensed keeps sending after the license lapses.

Inputs ​

The block has two modes. Simple (the default) is a form. Raw gives your JavaScript the broker's own client.

Simple mode ​

FieldTabRequiredDefaultWhat it does
IntegrationGeneralYesnoneA Kafka, NATS, SQS or RabbitMQ message queue integration, or a Redis integration to send to a stream.
DestinationGeneralYesnoneWhere the message goes inside that integration. See the table below. Supports js: expressions. In a list, a message can name its own.
Use Input as PayloadGeneralNooffSend the previous block's output as it is. A list sends one message per item; anything else is one message. The Message tab is hidden.
MessagesMessageNoSingleSingle sends one message. Bulk sends a list, one message per item.
PayloadMessageYes{}The message, as JSON (values support js:) or as JavaScript code that returns it.
KeyOptionsNononeKafka only. Messages with the same key go to the same partition, in order.
HeadersOptionsNononeKafka, NATS and RabbitMQ headers, or SQS message attributes. Redis streams have none.
Go to Error WhenOptionsNoAll messages failFor a list only. Decides when the Failure path runs. See Outputs.
Save output to variableNooffStore the output in outputs.<name>, on either path.

What to put in Destination:

BrokerDestination isExample
KafkaA topic nameorders
NATSA subject that one of your JetStream streams listens onorders.created
SQSThe queue's full URLhttps://sqs.us-east-1.amazonaws.com/123456789012/orders
RabbitMQThe routing key. With no Exchange set, this is the name of the queue to send to.orders
RedisThe stream's key. The stream is created on the first message.orders

Sending a list as one message

With Use Input as Payload on, a list always becomes one message per item. To send a whole list as one message, turn it off, keep Messages on Single, and set the payload to JavaScript: return input.

Broker options, also on the Options tab. The Name is what to use when a message in a list sets its own (see Lists of messages).

BrokerOptionNameWhat it does
KafkaPartitionpartitionSend to this partition. Blank lets Kafka pick from the key.
KafkaTimestamptimestampMilliseconds since 1970, or a date. Blank is now.
NATSMessage IDmsgIdJetStream drops a second message with the same ID inside the stream's duplicate window.
SQSDelay (seconds)delaySecondsHide the message for 0 to 900 seconds. Not for FIFO queues.
SQSMessage Group IDgroupIdRequired by FIFO queues. Messages in one group are delivered in order.
SQSDeduplication IDdeduplicationIdFIFO queues without content-based deduplication need one.
RabbitMQExchangeexchangeThe exchange to publish to. Blank sends straight to the queue named in Destination.
RabbitMQMessage IDmessageIdBlank gets a new unique ID. A RabbitMQ trigger builds meta.id from it.
RabbitMQContent TypecontentTypeBlank is application/json, or text/plain for text. Only the label changes, not the body.
RabbitMQExpiration (ms)expirationRabbitMQ drops the message if no one reads it in this many milliseconds. Blank keeps it.
RedisMax LengthmaxLenTrim the stream to about this many entries. Blank keeps everything.

Raw mode ​

Use Raw mode when the form can't do what you need, such as a broker feature it has no field for.

FieldTabRequiredDefaultWhat it does
IntegrationGeneralYesnoneAs in Simple mode.
CodeCodeYesnoneJavaScript with a client object, the broker's own client, already connected. It must return the result. The editor knows the client's types, so you get autocomplete.
Brokerclient is
KafkaA @platformatic/kafka producer with text keys and values, waiting for all in-sync replicas.
NATSA JetStream client from @nats-io/jetstream.
SQSThe SQS client from @aws-sdk/client-sqs, with methods like sendMessage.
RabbitMQA confirm channel from amqplib. Pass a callback to publish to wait for RabbitMQ's confirm.
RedisAn ioredis client.

Outputs ​

HandleRuns whenInput to the next block
Success (right)The broker accepted the message. For a list, see below.What the broker reported.
Failure (right, red)The send failed.{ error }, or for a list the full report.

What the broker reports for one message:

BrokerOutput
Kafka{ topic, partition, offset }
NATS{ stream, seq, duplicate }
SQS{ queueUrl, messageId, sequenceNumber } (sequence number on FIFO queues only)
RabbitMQ{ exchange, routingKey, messageId } (exchange is "" when sending straight to a queue)
Redis{ stream, id }

For a list, the output lists every message by its position in your list:

json
{
  "sent": [{ "index": 0, "topic": "orders", "partition": 1, "offset": "8812" }],
  "failed": [{ "index": 1, "error": "Topic not found" }]
}

A list is not all-or-nothing: messages that went through stay sent. The Go to Error When option decides which path runs:

Go to Error WhenSuccess runsFailure runs
All messages fail (default)At least one message was sent, or the list was empty.Every message failed.
Any message failsEvery message was sent.Even one message failed.

If you leave Failure unconnected, a failure goes to the Error Handler instead, like any other block's error.

Example ​

Publish an order to Kafka after saving it, keyed by customer so one customer's orders stay in order.

  • Destination: orders
  • Use Input as Payload: on (the previous block returns the saved order)
  • Key: js:return input.customerId

The next block receives { "topic": "orders", "partition": 2, "offset": "8813" }.

Send a list of reminders to SQS, one of them to a different queue. Set Destination to the main queue's URL and turn on Use Input as Payload. The previous block returns:

json
[
  { "userId": 1 },
  { "userId": 2 },
  { "payload": { "userId": 3 }, "destination": "https://sqs.us-east-1.amazonaws.com/123456789012/vip" }
]

Users 1 and 2 go to the main queue, user 3 to the vip queue. The next block receives a report like the one under Outputs.

How it behaves ​

What gets sent ​

PayloadSent as
TextThe text, unchanged.
Number or true/falseIts text, e.g. 42. A Fluxify trigger reads it back as a number or true/false.
Object or listJSON.
BigInt (large whole numbers from a database)Its digits as a JSON string, e.g. "9007199254740993", so no digits are lost.
Function, symbol, undefined, or an object that contains itselfNot sent. The message fails with the reason.

Text that looks like JSON

A Fluxify trigger tries to read every message as JSON. So the text "42" or "true" comes back as the number 42 or true. If it must stay text, send an object like { "value": "42" }.

On Redis, a stream entry holds named fields, not one body. An object's keys become the fields, and each value is sent as above. Anything else is sent in one field called data. A Redis Streams trigger hands the fields back as text.

On RabbitMQ, every message is persistent and must reach a queue: one that no queue takes fails instead of disappearing. See Sending to RabbitMQ for every default.

Lists of messages ​

  • Each item in the list is one message, and the item is the payload.
  • To give one message its own settings, wrap it: { "payload": ..., "destination": ..., "key": ..., "headers": {...} }, plus any broker option by its Name from the table above, such as "partition" or "groupId". These override the block's settings for that message only. Headers are merged with the block's.
  • Every message gets its own result, so one bad message doesn't stop the others.
  • You don't need to split large lists. SQS takes at most 10 messages or 256 KB per request, so the block sends in chunks of that size for you.

Waiting and timeouts ​

  • The block always waits for the broker to confirm, and never retries on its own.
  • It waits up to the integration's Send timeout, 30 seconds by default. After that the send fails with "No answer from the broker". A message that timed out may still have arrived.
  • One connection per integration is shared by every run on a worker. Editing the integration reconnects. In Raw mode, client.close() and similar methods are refused for this reason.

Raw examples ​

Kafka:

javascript
const result = await client.send({
  messages: [{ topic: "orders", key: "42", value: JSON.stringify({ id: 42 }) }],
});
return result.offsets;

NATS:

javascript
const ack = await client.publish("orders.created", JSON.stringify({ id: 42 }));
return { stream: ack.stream, seq: ack.seq };

SQS:

javascript
const sent = await client.sendMessage({
  QueueUrl: "https://sqs.us-east-1.amazonaws.com/123456789012/orders",
  MessageBody: JSON.stringify({ id: 42 }),
});
return sent.MessageId;

RabbitMQ:

javascript
await new Promise((resolve, reject) =>
  client.publish("", "orders", Buffer.from(JSON.stringify({ id: 42 })), { persistent: true },
    (error) => (error ? reject(error) : resolve())),
);
return "sent";

Redis:

javascript
return await client.xadd("orders", "*", "id", "42", "status", "new");

Released under the Apache License 2.0. Enterprise features are under the Fluxify Enterprise Edition License.