SDK of Diaspora Event Fabric: Resilience-enabling services for science from HPC to edge
Project description
Diaspora Event Fabric: Resilience-enabling services for science from HPC to edge
- Installation Instructions
- Use Diaspora Event SDK
Installation Instructions
Recommended Installation with Kafka Client Library: kafka-python
If you plan to utilize the KafkaAdmin
, KafkaProducer
, and KafkaConsumer
classes in the SDK, which are extensions of the respective classes from the kafka-python
library, we recommend installing the SDK with kafka-python
support. This is especially convenient for tutorial purposes and integrating Kafka functionalities in your projects with out-of-box configurations.
To install Diaspora Event SDK with kafka-python
, run:
pip install "diaspora-event-sdk[kafka-python]"
Installation Without Kafka Client Library
For scenarios where kafka-python
is not required or if you are using other client libraries to communicate with Kafka, you can install the SDK without this dependency.
To install the SDK without Kafka support, simply run:
pip install diaspora-event-sdk
Note that this option does not install the necessary dependency for KafkaAdmin
, KafkaProducer
, and KafkaConsumer
below to work.
Use Diaspora Event SDK
Use SDK to communicate with Kafka (kafka-python Required)
Register Topic (create topic ACLs)
Before you can create, describe, and delete topics we need to set the appropriate ACLs in ZooKeeper. Here we use the Client to register ACLs for the desired topic name.
from diaspora_event_sdk import Client as GlobusClient
c = GlobusClient()
topic = "topic-" + c.subject_openid[-12:]
print(c.register_topic(topic))
print(c.list_topics())
Create Topic
Now use the KafkaAdmin to create the topic.
from diaspora_event_sdk import KafkaAdmin, NewTopic
admin = KafkaAdmin()
print(admin.create_topics(new_topics=[
NewTopic(name=topic, num_partitions=1, replication_factor=1)]))
Start Producer
Once the topic is created we can publish to it. The KafkaProducer wraps the Python KafkaProducer Event publication can be either synchronous or asynchronous. Below demonstrates the synchronous approach.
from diaspora_event_sdk import KafkaProducer
producer = KafkaProducer()
future = producer.send(
topic, {'message': 'Synchronous message from Diaspora SDK'})
print(future.get(timeout=10))
Start Consumer
A consumer can be configured to monitor the topic and act on events as they are published. The KafkaConsumer wraps the Python KafkaConsumer. Here we use the auto_offset_reset
to consume from the first event published to the topic. Removing this field will have the consumer act only on new events.
from diaspora_event_sdk import KafkaConsumer
consumer = KafkaConsumer(topic, auto_offset_reset='earliest')
for msg in consumer:
print(msg)
Delete Topic
from diaspora_event_sdk import KafkaAdmin
admin = KafkaAdmin()
print(admin.delete_topics(topics=[topic]))
Unregister Topic (remove topic ACLs)
from diaspora_event_sdk import Client as GlobusClient
c = GlobusClient()
topic = "topic-" + c.subject_openid[-12:]
print(c.unregister_topic(topic))
print(c.list_topics())
Communicating with Kafka Using Your Preferred Client Library
Register and Unregister Topic
The steps are the same as above by using the register_topic
, unregister_topic
, and list_topics
methods from the Client
class.
Cluster Connection Details
Configuration | Value |
---|---|
Bootstrap Servers | MSK_SCRAM_ENDPOINT |
Security Protocol | SASL_SSL |
Sasl Mechanism | SCRAM-SHA-512 |
Api Version | 3.5.1 |
Username | (See instructions below) |
Password | (See instructions below) |
Execute the code snippet below to obtain your unique username and password for the Kafka cluster:
from diaspora_event_sdk import Client as GlobusClient
c = GlobusClient()
print(c.retrieve_key())
Advanced Usage
Password Refresh
In case that you need to invalidate all previously issued passwords and generate a new one, call the create_key
method from the Client
class
from diaspora_event_sdk import Client as GlobusClient
c = GlobusClient()
print(c.create_key())
Subsequent calls to retrieve_key
will return the new password from the cache. This cache is reset with a logout or a new create_key
call.
Common Issues
ImportError: cannot import name 'KafkaAdmin' from 'diaspora_event_sdk'
It seems that you ran pip install diaspora-event-sdk
to install the Diaspora Event SDK without kafka-python
. Run pip install kafka-python
to install the necessary dependency for our KafkaAdmin
, KafkaProducer
, and KafkaConsumer
classes.
kafka.errors.NoBrokersAvailable and kafka.errors.NodeNotReadyError
These messages might pop up if create_key
is called shortly before instanciating a Kafka client. This is because there's a delay for AWS Secret Manager to associate the newly generated credential with MSK. Note that create_key
is called the first time you create one of these clients. Please wait a while (around 1 minute) and retry.
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Hashes for diaspora-event-sdk-0.0.10.tar.gz
Algorithm | Hash digest | |
---|---|---|
SHA256 | 66a00102ce1e7af772c336ef7bb4380c5d062a4d30b7e2ebfb4ba48ea5413d83 |
|
MD5 | fe5e0cfc8b84358af5ae4215e39b9a9d |
|
BLAKE2b-256 | be62fa0e071ab343c3535cb33ef25e591b60b289700a5c442129c2140c2bbce8 |
Hashes for diaspora_event_sdk-0.0.10-py3-none-any.whl
Algorithm | Hash digest | |
---|---|---|
SHA256 | 7fc7a0075ffeea383e1de8672ad59cb30b924acff5aab7180edcfb6f3c6c42b1 |
|
MD5 | f098e7817b0887e8e9296f7785b01a69 |
|
BLAKE2b-256 | eb0c0b8cb1a6cb34e385d2a645df1e6d774b38484d83f8338ba4da75cc0f26de |