Event Brokers
Event Brokers
An event broker allows you to connect your running assistant to other services that process the data coming in from conversations. For example, you could connect your live assistant to Rasa X to review and annotate conversations or forward messages to an external analytics service. The event broker publishes messages to a message streaming service, also known as a message broker, to forward Rasa Events from the Rasa server to other services.
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).
Pika Event Broker
The example implementation we’re going to show you here uses Pika, the Python client library for RabbitMQ.
Adding a Pika Event Broker Using the Endpoint Configuration
You can instruct Rasa to stream all events to your Pika event broker by adding an event_broker section to your endpoints.yml:
event_broker:
type: pika
url: localhost
username: username
password: password
queues:
- queue-1
# you may supply more than one queue to publish to
# - queue-2
# - queue-3
Rasa will automatically start streaming events when you restart the Rasa server.
Adding a Pika Event Broker in Python
Here is how you add it using Python code:
from rasa.core.brokers.pika import PikaEventBroker
from rasa.core.tracker_store import InMemoryTrackerStore
pika_broker = PikaEventBroker('localhost',
'username',
'password',
queues=['rasa_events'])
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 needs to implement Pika’s start_consuming() method with a callback action. Here’s a simple example:
import json
import pika
def _callback(self, 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(_callback,
queue='rasa_events',
no_ack=True)
channel.start_consuming()
Kafka Event Broker
It is possible to use Kafka as the main broker for your events. In this example, we are going to use the python-kafka library, a Kafka client written in Python.
Adding a Kafka Event Broker Using the Endpoint Configuration
You can instruct Rasa to stream all events to your Kafka event broker by adding an event_broker section to your endpoints.yml:
Using SASL_PLAINTEXT protocol the endpoints file must have the following entries:
event_broker:
url: localhost
partition_by_sender: True
sasl_username: username
sasl_password: password
sasl_mechanism: PLAIN
topic: topic
security_protocol: SASL_PLAINTEXT
type: kafka
If using SSL protocol, the endpoints file should look like:
event_broker:
url: localhost
topic: topic
security_protocol: SSL
ssl_cafile: CARoot.pem
ssl_certfile: certificate.pem
ssl_keyfile: key.pem
ssl_check_hostname: True
type: kafka
Adding a Kafka Broker in Python
The code below shows an example on how to instantiate a Kafka producer in your script:
from rasa.core.brokers.kafka import KafkaEventBroker
from rasa.core.tracker_store import InMemoryTrackerStore
kafka_broker = KafkaEventBroker(host='localhost:9092',
topic='rasa_events')
tracker_store = InMemoryTrackerStore(domain=domain, event_broker=kafka_broker)
Partition Key
Rasa Open Source’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.
event_broker:
type: kafka
partition_by_sender: True
security_protocol: PLAINTEXT
topic: topic
url: localhost
client_id: kafka-python-rasa
Authentication and Authorization
Rasa’s Kafka producer accepts two types of security protocols - SASL_PLAINTEXT and SSL.
For development environment, you can use simple authentication with SASL_PLAINTEXT.
kafka_broker = KafkaEventBroker(host='kafka_broker:9092',
sasl_plain_username='kafka_username',
sasl_plain_password='kafka_password',
security_protocol='SASL_PLAINTEXT',
topic='rasa_events')
If the clients or the brokers in the Kafka cluster are located in different machines, it’s important to use SSL protocol to assure encryption of data and client authentication.
kafka_broker = KafkaEventBroker(host='kafka_broker:9092',
ssl_cafile='CARoot.pem',
ssl_certfile='certificate.pem',
ssl_keyfile='key.pem',
ssl_check_hostname=True,
security_protocol='SSL',
topic='rasa_events')
Implementing a Kafka Event Consumer
The parameters used to create a Kafka consumer are the same used in the producer creation.
from kafka import KafkaConsumer
from json import loads
consumer = KafkaConsumer('rasa_events',
bootstrap_servers=['localhost:29093'],
value_deserializer=lambda m: json.loads(m.decode('utf-8')),
security_protocol='SSL',
ssl_check_hostname=False,
ssl_cafile='CARoot.pem',
ssl_certfile='certificate.pem',
ssl_keyfile='key.pem')
for message in consumer:
print(message.value)
SQL Event Broker
It is possible to use an SQL database as an event broker. Connections to databases are established using SQLAlchemy.
Adding a SQL Event Broker Using the Endpoint Configuration
To instruct Rasa to save all events to your SQL event broker, add an event_broker section to your endpoints.yml. For example:
event_broker:
type: SQL
dialect: sqlite
db: events.db
PostgreSQL databases can also be utilized:
event_broker:
type: SQL
host: 127.0.0.1
port: 5432
dialect: postgresql
username: myuser
password: mypassword
db: mydatabase
With this configuration applied, Rasa will create a table called events on the database.