Skip to content

RabbitMQ Source Connector ​

stream sourcescan table source

The RabbitMQ source connector reads messages from an AMQP 0-9-1 queue or exchange.

You can use the RabbitMQ connector as a source connector and as a sink connector.

Configuration Properties ​

Configure the connector in the WITH clause of a CREATE STREAM statement, or configure the connector in a configuration file.

PropertyTypeDefaultDescription
server (or url)string"amqp://guest:[email protected]:5672/%2f"The AMQP broker connection URL. Supports TLS URLs (amqps://) and secret templates.
queuestring""The name of the queue to consume.
exchangestring""The exchange name to bind the queue to.
routingKeystring""The routing key or routing pattern for queue binding.
durablebooleantrueWhen true, the queue survives broker restarts.
autoDeletebooleanfalseWhen true, the broker removes the queue after consumers disconnect.
exclusivebooleanfalseWhen true, only this connection can access the queue.
prefetchCountinteger100The AMQP basic QoS prefetch count for message flow control.

Stream Definition Examples ​

Connect to a Standard Queue ​

The following statement creates a stream that consumes JSON data from a queue named telemetry_queue:

sql
CREATE STREAM rabbit_telemetry () WITH (
    TYPE = "rabbitmq",
    SERVER = "amqp://guest:[email protected]:5672/%2f",
    QUEUE = "telemetry_queue",
    FORMAT = "JSON"
);

Connect with Topic Exchange Binding ​

The following statement binds a queue to a topic exchange with a routing key pattern:

sql
CREATE STREAM factory_sensors (
    device_id STRING,
    temperature FLOAT,
    humidity FLOAT
) WITH (
    TYPE = "rabbitmq",
    SERVER = "amqp://guest:[email protected]:5672/%2f",
    QUEUE = "sensor_events",
    EXCHANGE = "amq.topic",
    ROUTINGKEY = "sensors.temperature.*",
    DURABLE = "true",
    PREFETCHCOUNT = "50",
    FORMAT = "JSON"
);

Security Credentials ​

Do not store plain-text passwords in stream definitions.

You can retrieve passwords from environment variables or HashiCorp Vault.

Example with Environment Variable ​

sql
CREATE STREAM secure_stream () WITH (
    TYPE = "rabbitmq",
    SERVER = "amqp://app_user:{{env://RABBIT_PASSWORD}}@rabbitmq.internal:5672/%2f",
    QUEUE = "telemetry",
    FORMAT = "JSON"
);

Example with HashiCorp Vault ​

sql
CREATE STREAM vault_stream () WITH (
    TYPE = "rabbitmq",
    SERVER = "amqp://app_user:{{vault://secret/rabbitmq#password}}@rabbitmq.internal:5672/%2f",
    QUEUE = "high_security_events",
    FORMAT = "JSON"
);

Refer to the Dynamic Secrets documentation for more details.

Rule Example ​

The following rule calculates the average temperature for each device in a 10-second tumbling window:

sql
SELECT 
    device_id, 
    avg(temperature) AS avg_temp,
    count(*) AS alert_count
FROM factory_sensors
WHERE temperature > 75.0
GROUP BY device_id, TumblingWindow(ss, 10);

Released under the Apache-2.0 / MIT License.