Skip to main content

PyPI

Message Bus Schema and Message Definition

Basic framework for a Fabric Message Bus Schema and Messages for Inter Actor Communication

Overview

Fabric communication across various actors in Control and Measurement Framework is implemented using Apache Kafka.

Apache Kafka is a distributed system designed for streams. It is built to be fault-tolerant, high-throughput, horizontally scalable, and allows geographically distributed data streams and stream processing applications.

Kafka enables event driven implementation of various actors/services. Events are both a Fact and a Trigger. Each fabric actor will be a producer for one topic following the Single Writer Principle and would subscribe to topics from other actors for communication. Messages are exchanged over Kafka using Apache Avro data serialization system.

Requirements

  • Python 3.7+
  • confluent-kafka
  • confluent-kafka[avro]

Installation

$ pip3 install .

Usage

This package implements the interface for producer/consumer APIs to push/read messages to/from Kafka via Avro serialization.

Message and Schema

User is expected to inherit IMessage class(message.py) to define it's own members and over ride to_dict() function. It is also required to define the corresponding AVRO schema pertaining to the derived class. This new schema shall be used in producer and consumers.

Example schema for basic IMessage class is available in (schema/message.avsc)

Producers

AvroProducerApi class implements the base functionality for an Avro Kafka producer. User is expected to inherit this class and override delivery_report method to handle message delivery for asynchronous produce.

Example for usage available at the end of producer.py

Consumers

AvroConsumerApi class implements the base functionality for an Avro Kafka consumer. User is expected to inherit this class and override process_message method to handle message processing for incoming message.

Example for usage available at the end of consumer.py

Admin API

AdminApi class provides support to carry out basic admin functions like create/delete topics/partions etc.

How to bring up a test Kafka cluster to test

Generate Credentials

You must generate CA certificates (or use yours if you already have one) and then generate a keystore and truststore for brokers and clients.

cd $(pwd)/secrets
./create-certs.sh
(Type yes for all "Trust this certificate? [no]:" prompts.)
cd -

Set the environment variable for the secrets directory. This is used in later commands. Make sure that you are in the MessageBus directory.

export KAFKA_SSL_SECRETS_DIR=$(pwd)/secrets

Bring up the containers

You can use the docker-compose.yaml file to bring up a simple Kafka cluster containing

  • broker
  • zookeeper
  • schema registry

Use the below command to bring up the cluster

docker-compose up -d

This should bring up following containers:

docker ps
CONTAINER ID        IMAGE                                    COMMAND                  CREATED             STATUS              PORTS                                                                                        NAMES
189ba0e70b97        confluentinc/cp-schema-registry:latest   "/etc/confluent/dock…"   58 seconds ago      Up 58 seconds       0.0.0.0:8081->8081/tcp                                                                       schemaregistry
49616f1c9b0a        confluentinc/cp-kafka:latest             "/etc/confluent/dock…"   59 seconds ago      Up 58 seconds       0.0.0.0:9092->9092/tcp, 0.0.0.0:19092->19092/tcp                                             broker1
c9d19c82558d        confluentinc/cp-zookeeper:latest         "/etc/confluent/dock…"   59 seconds ago      Up 59 seconds       2888/tcp, 0.0.0.0:2181->2181/tcp, 3888/tcp                                                   zookeeper

Release files for fabric_message_bus 2.0.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for fabric_message_bus 2.0.1
File Size Uploaded
fabric_message_bus-2.0.1.tar.gz 44.6 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for fabric_message_bus 2.0.1
File Interpreter ABI Platform
fabric_message_bus-2.0.1-py3-none-any.whl Python 3 none any Details

Total release size: 190.1 kB

Release files / fabric_message_bus-2.0.1.tar.gz

Download URL fabric_message_bus-2.0.1.tar.gz
Size 44.6 kB
Tags Source
SHA-256 checksum
How to use checksums
a2ff7febca20c8b26951369be6a64f149821290ee47c16e824c2054cf6466bbb
BLAKE2b-256 checksum
How to use checksums
d31022a111319249d47224710750e2d23696b4769c7216176451b37e149f1d75
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via python-requests/2.32.5

Release files / fabric_message_bus-2.0.1-py3-none-any.whl

Download URL fabric_message_bus-2.0.1-py3-none-any.whl
Size 145.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
a28222ddaf62e0667ee68923c5cdc20cc2f4403bd19d4731101dc9b96b6cc350
BLAKE2b-256 checksum
How to use checksums
02ff4237104cb057b715842a8e02307c93106671f97fe86b5798f3f6d2b106c6
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via python-requests/2.32.5

Release history Release notifications | RSS feed

This release

2.0.1 This release

2 release files

2.0.0

2 release files

1.10.1

2 release files

1.10.0

2 release files

1.9.0

2 release files

1.7.0

2 release files

1.6.2

2 release files

1.6.1

2 release files

1.6.0

2 release files

1.5.0

2 release files

1.4.0

2 release files

1.3.0

2 release files

1.2.3

2 release files

1.2.2

2 release files

1.2.1

2 release files

1.2

2 release files

1.0.6

2 release files

1.0.5

2 release files

1.0.4

2 release files

1.0.3

2 release files

1.0.2

2 release files

1.0.1

2 release files

1.0

2 release files

0.14

2 release files

0.13

2 release files

0.12

2 release files

0.11

2 release files

0.10

2 release files

0.9

2 release files

0.8

2 release files

0.7

2 release files

0.6

2 release files

0.5

2 release files

0.4

2 release files

0.3

2 release files

0.2

2 release files

0.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page