Event Brokers | Rasa Documentation
Format
All events are streamed to the broker as serialized dictionaries every time the tracker updates its state. An example event emitted from the default tracker looks like this:
{
"sender_id": "default",
"timestamp": 1528402837.617099,
"event": "bot",
"text": "what your bot said",
"data": "some data about e.g. attachments",
"metadata": {
"a key": "a value"
}
}
The event field takes the event's type_name (for more on event types, check out the events docs).
Kafka Event Broker
Kafka is recommended for all assistants at scale. Kafka is a requirement when streaming events to Rasa Pro Services or Rasa Studio.
Rasa uses the confluent-kafka library, a Kafka client written in Python.
Configuration
To use Kafka as an event broker in Rasa, you need to set it as an event_broker in your endpoints.yml file. For Kafka, Rasa supports following properties in event_broker section:
endpoints.yml
event_broker:
type: kafka
url: localhost:9092 # required, url to your kafka broker
# Optional properties
topic: rasa_core_events # topic to which events are published
security_protocol: "SASL_PLAINTEXT" # security protocol to use, available options are: PLAINTEXT, SASL_PLAINTEXT, SSL, SASL_SSL
partition_by_sender: False # should events be partitioned by sender id
client_id: # ID to use for the producer of the events
# SASL configuration, optional
sasl_mechanism: "PLAIN" # SASL mechanism to use, available options are: PLAIN, GSSAPI, OAUTHBEARER, SCRAM-SHA-256, SCRAM-SHA-512
sasl_username: # username to use for authentication, only if security_protocol is SASL_PLAINTEXT or SASL_SSL
sasl_password: # password to use for authentication, only if security_protocol is SASL_PLAINTEXT or SASL_SSL
# TLS configuration, optional
ssl_cafile: # path to the CA certificate file, only if security_protocol is SSL or SASL_SSL
ssl_certfile: # path to the client certificate file, only if security_protocol is SSL or SASL_SSL and Kafka is configured to use client authentication
ssl_keyfile: # path to the client key file, only if security_protocol is SSL or SASL_SSL and Kafka is configured to use client authentication
ssl_check_hostname: False # whether to check the hostname of the broker against the certificate, default is True
# PII management configuration, optional
stream_pii: False # whether to stream PII events, default is True
anonymization_topics: # list of topics to publish anonymized events to, default is []
- anonymized_topic_1
Partition Key
Rasa's Kafka producer can optionally be configured to partition messages by conversation ID. This can be configured by setting partition_by_sender in the endpoints.yml file to True. By default, this parameter is set to False and the producer will randomly assign a partition to each message.
Authentication and Authorization
Rasa's Kafka producer accepts the following types of security protocols: SASL_PLAINTEXT, SSL, PLAINTEXT and SASL_SSL.
For development environments, or if the brokers servers and clients are located into the same machine, you can use simple authentication with SASL_PLAINTEXT or PLAINTEXT. By using this protocol, the credentials and messages exchanged between the clients and servers will be sent in plaintext. Thus, this is not the most secure approach, but since it's simple to configure, it is useful for simple cluster configurations. SASL_PLAINTEXT protocol requires the setup of the username and password previously configured in the broker server.
If the clients or the brokers in the kafka cluster are located in different machines, it's important to use the SSL or SASL_SSL protocol to ensure encryption of data and client authentication. After generating valid certificates for the brokers and the clients, the path to the certificate and key generated for the producer must be provided as arguments, as well as the CA's root certificate.
If using the GSSAPI SASL mechanism, you will need to additionally install python-gssapi and the necessary C library Kerberos dependencies.
Using IAM roles to authenticate to AWS Managed Streaming for Apache Kafka (MSK)
You can use IAM authentication to connect to AWS MSK without needing to provide a username and password.
If your Rasa instance is running on an AWS service that supports IAM roles (e.g. EC2), you can use IAM authentication to connect to an AWS MSK cluster without needing to provide a username and password.
To do so, you need to ensure that your MSK cluster is configured to allow IAM authentication. Ensure that the topic(s) you want to use are created on the MSK cluster. You also need to set up your Rasa instance with an appropriate IAM role that has permissions to access the MSK cluster.
The IAM role should include the following permissions:
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"kafka-cluster:Connect",
"kafka-cluster:DescribeCluster"
],
"Resource": [
"arn:aws:kafka:<region_name>:<account_id>:cluster/<cluster_name>/<cluster_uuid>"
]
},
{
"Effect": "Allow",
"Action": [
"kafka-cluster:*Topic*",
"kafka-cluster:WriteData",
"kafka-cluster:ReadData"
],
"Resource": [
"arn:aws:kafka:<region_name>:<account_id>:topic/<cluster_name>/*"
]
},
{
"Effect": "Allow",
"Action": [
"kafka-cluster:AlterGroup",
"kafka-cluster:DescribeGroup"
],
"Resource": [
"arn:aws:kafka:<region_name>:<account_id>:group/<cluster_name>/*"
]
}
]
}
Once you have set up the IAM role and configured your AWS MSK cluster, you can configure the following environment variables in your Rasa instance:
IAM_CLOUD_PROVIDER: Set this toaws.AWS_DEFAULT_REGION: Set this to the AWS region where your MSK cluster is located.KAFKA_MSK_AWS_IAM_ENABLED: Set this totrueto enable IAM authentication for MSK connections, otherwise leave it unset or set tofalse.
You also need to configure the Kafka event broker in your endpoints.yml file to use the SASL_SSL security protocol and OAUTHBEARER SASL mechanism, as well as provide the path to the root certificate from Amazon Trust Services in the ssl_cafile property.
endpoints.yml
event_broker:
type: kafka
url: <your-msk-broker-url>
topic: rasa_core_events
security_protocol: "SASL_SSL"
sasl_mechanism: "OAUTHBEARER"
ssl_cafile: <path-to-amazon-trust-services-root-certificate>
partition_by_sender: true
client_id: <your-client-id>
ssl_check_hostname: true
When you start your Rasa instance, it will use the IAM role to generate temporary credentials to log in to the AWS MSK cluster instead of using static credentials. The temporary credentials will be automatically refreshed every 15 minutes.
Example Configurations
Without authentication
To set up Rasa with Kafka which does not require authentication nor TLS handshake, use the following config as an example:
endpoints.yml
event_broker:
type: kafka
security_protocol: PLAINTEXT
topic: topic
url: localhost
To set up Rasa with Kafka which does not require authentication but uses TLS handshake, use the following config as an example:
endpoints.yml
event_broker:
type: kafka
security_protocol: SSL
topic: topic
url: localhost
ssl_cafile: CARoot.pem
ssl_certfile: certificate.pem
ssl_keyfile: key.pem
ssl_check_hostname: True
With authentication
To set up Rasa with Kafka which requires authentication but does not use TLS handshake, use the following config as an example:
endpoints.yml
event_broker:
type: kafka
security_protocol: SASL_PLAINTEXT
topic: topic
url: localhost
sasl_username: username
sasl_password: password
sasl_mechanism: PLAIN
To set up Rasa with Kafka which requires authentication and uses TLS handshake, use the following config as an example:
endpoints.yml
event_broker:
type: kafka
security_protocol: SASL_SSL
topic: topic
url: localhost
sasl_username: username
sasl_password: password
sasl_mechanism: PLAIN
ssl_cafile: CARoot.pem
ssl_certfile: certificate.pem
ssl_keyfile: key.pem
ssl_check_hostname: True
Make sure that the SASL mechanism is set according to the broker configuration. You can also use GSSAPI, OAUTHBEARER, SCRAM-SHA-256 or SCRAM-SHA-512 if your broker is configured to use it for the exposed URL endpoint.
Sending Events to Multiple Queues
Kafka does not allow you to configure multiple topics. However, multiple consumers can read from the same queue as long as they are in different consumer groups. Each consumer group will process all events independent of each other.
Disabling Publishing of Un-anonymised Events
You can configure the event broker to not publish un-anonymised events to the configured topic. This is done by setting the stream_pii parameter in the endpoints.yml file to false.
Sending Anonymized Events
If you have the PII management capability enabled, you can configure the event broker to publish anonymised events to a different topic. This is done by setting the anonymization_topics parameter in the endpoints.yml file to a list of topics.
Non-Blocking Publishing
You can configure Kafka to publish events without blocking the event loop by setting type: concurrent_kafka instead of type: kafka. By default, type: kafka publishes events synchronously on the event loop. Each publish call waits for the Kafka broker to acknowledge the event before returning. For latency-sensitive deployments, use type: concurrent_kafka instead.
endpoints.yml
event_broker:
type: concurrent_kafka
url: localhost:9092
topic: rasa_core_events
# All options from type: kafka are supported
security_protocol: "SASL_PLAINTEXT"
sasl_mechanism: "PLAIN"
sasl_username: myuser
sasl_password: mypassword
partition_by_sender: True
# Additional option: number of background publish threads (default: 1)
executor_max_workers: 1
Ordering with multiple workers
When executor_max_workers is greater than 1, concurrent retries can reorder events for the same sender. Set partition_by_sender: True to ensure events for a given conversation always land on the same Kafka partition, but note that retry races between threads may still affect ordering within a partition. For strict ordering guarantees, keep executor_max_workers: 1 (the default).
Pika Event Broker for RabbitMQ
Rasa uses Pika, the Python client library for RabbitMQ.
Configuration
To use RabbitMQ as an event broker in Rasa, you need to set it as an event_broker in your endpoints.yml file. For RabbitMQ, Rasa supports following properties in even_broker section:
endpoints.yml
event_broker:
type: pika
host: # required, hostname of your RabbitMQ broker e.g. localhost
port: 5672 # port of your RabbitMQ broker e.g. 5672
exchange_name: "rasa-exchange" # exchange name to use for publishing events
username: # required, username to use for authentication
password: # required, password to use for authentication
connection_attempts: 20 # number of connection attempts to make before giving up
retry_delay_in_seconds: 5 # time to wait before retrying connection attempts
raise_on_failure: False #whether to raise an exception on connection failure
should_keep_unpublished_messages: True # whether to keep unpublished messages in memory
queues: # list of queues to publish events to"
- rasa_core_events # default
stream_pii: False # whether to publish un-anonymised events to the configured queues. defaults to True
anonymization_queues: # list of queues to publish anonymized events to. defaults to []
- anonymized_event_queue
Additionally, you can use TLS with RabbitMQ by setting the following environment variables:
RABBITMQ_SSL_CLIENT_CERTIFICATE: path to the SSL client certificateRABBITMQ_SSL_CLIENT_KEY: path to the SSL client key
Example Configurations
To set up Rasa with Pika for RabbitMQ use the following config as an example:
endpoints.yml
event_broker:
type: pika
url: localhost
username: username
password: password
queues:
- queue-1
exchange_name: exchange
Adding a Pika Event Broker in Python
Here is how you add it using Python code:
import asyncio
from rasa.core.brokers.pika import PikaEventBroker
from rasa.core.tracker_store import InMemoryTrackerStore
pika_broker = PikaEventBroker('localhost',
'username',
'password',
queues=['rasa_events'],
event_loop=event_loop
)
asyncio.run(pika_broker.connect())
tracker_store = InMemoryTrackerStore(domain=domain, event_broker=pika_broker)
Implementing a Pika Event Consumer
You need to have a RabbitMQ server running, as well as another application that consumes the events. This consumer to needs to implement Pika's start_consuming() method with a callback action. Here's a simple example:
import json
import pika
def _callback(ch, method, properties, body):
print("Received event {}".format(json.loads(body)))
if __name__ == "__main__":
credentials = pika.PlainCredentials("username", "password")
connection = pika.BlockingConnection(
pika.ConnectionParameters("rabbit", credentials=credentials)
)
channel = connection.channel()
channel.basic_consume(queue="rasa_events", on_message_callback=_callback, auto_ack=True)
channel.start_consuming()
Sending Events to Multiple Queues
You can specify multiple event queues to publish events to. This should work for all event brokers supported by Pika (e.g. RabbitMQ).
Disabling Publishing of Un-anonymised Events
By default, Rasa will publish un-anonymised events to the configured queues. If you want to disable this and only publish anonymized events, set stream_pii: false in your event_broker configuration.
Publishing Anonymized Events
If you want to publish anonymized events to a different queue, you can set the anonymization_queues property in your event_broker configuration.
SQL Event Broker
It is possible to use an SQL database as an event broker. Connections to databases are established using SQLAlchemy, a Python library which can interact with many different types of SQL databases.
To set up Rasa with SQL event broker the following steps are required:
- Add required configuration to your
endpoints.yml
When using SQLite:
endpoints.yml
event_broker:
type: SQL
dialect: sqlite
db: events.db
When using PostgreSQL:
endpoints.yml
event_broker:
type: SQL
url: 127.0.0.1
port: 5432
dialect: postgresql
username: myuser
password: mypassword
db: mydatabase
- To start the Rasa server using your SQL backend, add the
--endpointsflag.
rasa run -m models --endpoints endpoints.yml
FileEventBroker
It is possible to use the FileEventBroker as an event broker. This implementation will log events to a file in json format.
You can provide a path key in the endpoints.yml file if you wish to override the default file name: rasa_event.log.
Custom Event Broker
If you need an event broker which is not available out of the box, you can implement your own by extending the base class EventBroker.
To set up Rasa with your custom event broker the following steps are required:
- Add required configuration to your
endpoints.yml
endpoints.yml
event_broker:
type: path.to.your.module.Class
url: localhost
a_parameter: a value
another_parameter: another value
- To start the Rasa server using your custom backend, add the
--endpointsflag.
rasa run -m models --endpoints endpoints.yml