Fabric events
Send events to an eventstream, run a service for each event Fabric raises and read recent events from an eventhouse.
This page will show you how to work with Microsoft Fabric events in Python:
- Sending an event to an eventstream
- Receiving events that Fabric raises
- Reading recent events from an eventhouse
Under the hood, Fabric receives events in an eventstream and stores them in an eventhouse:
- An eventstream exposes a custom endpoint, which is a Kafka endpoint, so services send events to it and receive events from it through Zato's Kafka connections
- An eventhouse answers KQL queries over HTTPS, so services read events from it through a Zato REST connection
Connections
The services on this page use three connections, created in the Zato Dashboard:
| Connection | Kind | Where in the Dashboard | Used for |
|---|---|---|---|
Fabric Events | Kafka outgoing connection | Connections → Message queues → Kafka → Outgoing connections | Sending events to the eventstream |
Fabric Alerts | Kafka channel | Connections → Message queues → Kafka → Channels | Running a service for each event the eventstream sends out |
Operations Events | REST outgoing connection | Connections → REST → Outgoing connections | Running KQL queries against the eventhouse |
Their values come from Fabric:
| Connection | In Fabric | Field | Value |
|---|---|---|---|
Fabric Events, Fabric Alerts | Eventstream → custom endpoint → Keys tab | Address | The bootstrap server |
Fabric Events, Fabric Alerts | Eventstream → custom endpoint → Keys tab | Topic | The topic name |
Operations Events | Eventhouse → overview page | Host | The Query URI |
Operations Events | Always the same | URL path | /v1/rest/query |
Each one signs in through a Bearer token security definition:
Fabric EventsandFabric Alertssign in as the app registration from the tutorial, which needs to be allowed to send to and receive from the eventstreamOperations Eventssigns in as an app registration of its own, created the same way as in the tutorial, because the eventhouse takes a scope of its own
Every field of each connection, with the exact scope its Bearer token needs, is in the events reference.
Sending an event
In this example, a service passes admissions from one system to another and each admission should also reach Fabric right away, not with the nightly load, so that a Real-Time Dashboard shows it within a second.
The service below will:
- Build the event from the admission it received
- Send it to the eventstream, which turns the dict into JSON
# -*- coding: utf-8 -*-
# Zato
from zato.server.service import Service
class PassAdmission(Service):
name = 'fabric.pass-admission'
input = 'admission_id', 'location', 'admitted_at'
def handle(self):
admission = self.request.input
# The event the eventstream receives ..
event = {
'event_type': 'admission',
'location': admission.location,
'occurred_at': admission.admitted_at,
'admission_id': admission.admission_id,
}
# .. get a Kafka connection ..
conn = self.out.kafka['Fabric Events']
# .. and send the event.
conn.send(event)
A second after invoking it, the event is in the eventstream's data preview in Fabric.
Receiving events
In this example, Fabric raises an event when the stock of an item falls below its reorder level, and the Fabric Alerts channel runs the service below for each such event.
The service below will:
- Read the event, which arrives as the JSON the eventstream sent
- Log its fields, which is where a real service would notify another system
# -*- coding: utf-8 -*-
# stdlib
import json
# Zato
from zato.server.service import Service
class NotifyPurchasing(Service):
name = 'fabric.notify-purchasing'
def handle(self):
# The event, as the eventstream sent it ..
event = json.loads(self.request.raw_request)
# .. its fields ..
item_id = event['item_id']
location = event['location']
quantity = event['quantity']
reorder_level = event['reorder_level']
# .. and log them.
self.logger.info(f'Low stock -> {item_id} at {location}')
self.logger.info(f'Left -> {quantity}')
self.logger.info(f'Reorder at -> {reorder_level}')
When Fabric raises the event, the server log has the line within the same second.
Reading recent events
In this example, the eventstream writes its events to the Events table of the Operations Events eventhouse, and a service needs the number of appointments cancelled in the last 15 minutes, per location.
The service below will:
- Run a KQL query against the eventhouse, over the
Operations EventsREST connection - Turn the reply, which is a list of columns and a list of rows, into a list of dicts
- Return the rows
# -*- coding: utf-8 -*-
# Zato
from zato.server.service import Service
class RecentCancellations(Service):
name = 'fabric.recent-cancellations'
def handle(self):
# The query to run ..
query = """
Events
| where event_type == 'appointment_cancelled'
| where occurred_at > ago(15m)
| summarize cancelled = count() by location
"""
# .. and the database to run it against ..
request = {
'db': 'Operations Events',
'csl': query,
}
# .. get a REST connection to the eventhouse ..
conn = self.rest['Operations Events']
# .. run the query ..
response = conn.post(self.cid, request)
# .. the result is the first table of the reply ..
tables = response.data['Tables']
result = tables[0]
# .. its column names are in one list ..
column_names = []
for column in result['Columns']:
column_names.append(column['ColumnName'])
# .. and the values of each row in another, in the same order ..
rows = []
for values in result['Rows']:
# .. so build a dict out of each name and its value ..
row = {}
for name, value in zip(column_names, values):
row[name] = value
rows.append(row)
# .. and return the rows to our caller.
self.response.payload = rows
After invoking the service you'll see:
See also
| Page | What it covers |
|---|---|
| Fabric API - Events | Every field of the Kafka and REST connections for eventstreams and eventhouses |
| Your first Fabric integration | The app registration the connections sign in as |