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:

ConnectionKindWhere in the DashboardUsed for
Fabric EventsKafka outgoing connectionConnections → Message queues → Kafka → Outgoing connectionsSending events to the eventstream
Fabric AlertsKafka channelConnections → Message queues → Kafka → ChannelsRunning a service for each event the eventstream sends out
Operations EventsREST outgoing connectionConnections → REST → Outgoing connectionsRunning KQL queries against the eventhouse

Their values come from Fabric:

ConnectionIn FabricFieldValue
Fabric Events, Fabric AlertsEventstream → custom endpoint → Keys tabAddressThe bootstrap server
Fabric Events, Fabric AlertsEventstream → custom endpoint → Keys tabTopicThe topic name
Operations EventsEventhouse → overview pageHostThe Query URI
Operations EventsAlways the sameURL path/v1/rest/query

Each one signs in through a Bearer token security definition:

  • Fabric Events and Fabric Alerts sign in as the app registration from the tutorial, which needs to be allowed to send to and receive from the eventstream
  • Operations Events signs 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 Events REST 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:

[{"location": "Maple Grove", "cancelled": 1}]

See also

PageWhat it covers
Fabric API - EventsEvery field of the Kafka and REST connections for eventstreams and eventhouses
Your first Fabric integrationThe app registration the connections sign in as

Learn more