Skip to main content

A library that enables building simple Kafka consumer/producer micro/nano services

Project description

kme

A Python library that enables building simple Kafka consumer/producer micro/nano services

This is a library that wraps the kafka-python library with some message routing helpers.

With this library, you can create micro/nano services (think like a Function, FaaS) that listens to one or more Kafka topics and processes the message, then based on configuration can send a new message.

I developed this with the intent to create a distributed workflow for various automation activities in my Kubernetes cluster and outside as well.

The workflow that I initially designed this for was:

  1. On a schedule, submit a message to the "wakup-computer" topic. This message contains details on which computer to wake up and how we determine when the computer is ready for work, and "completion topic", where we send a message when the computer is alive! In my case, I have an old server with a ton of hard drives, which I use for cold backups. I only need it on while backups are running.
  2. The completion topic, lets call it "run-backups" would accept the message and perform the backup processes. Once it has completed, it will send a new message to the "shutdown-computer" topic.
  3. The "shutdown-computer" topic has a consumer group for each computer which is included in this flow, and it receives the message and sees that it needs to shut down and does so.

I also invision using this workflow for my on-premise Kubernetes cluster so that I can "auto scale" nodes up/down when needed, reusing the "wakup-computer" and "shutdown-computer" topics and nano-services.

Usage

from kme import KMEMessage, KME
import os

# produce a message with a completion topic
def kme_producer():
    message = KMEMessage(topic=os.environ.get('KAFKA_TOPIC'))
    message.message = {'foo': 'bar', 'bar': 'foo'}
    message.completion_topic = 'foobar'
    k_client = KME(bootstrap_servers=[os.environ.get('KAFKA_BOOSTRAP_SERVER')])
    k_client.send_message(message=message)

# setup a consumer
def kme_consumer():
    k_client = KME(bootstrap_servers=[os.environ.get('KAFKA_BOOSTRAP_SERVER')])
    print(f"Subscribing to {os.environ.get('KAFKA_TOPIC')}", flush=True)
    # IMPORTANT - note the callback value!!
    k_client.subscribe(os.environ.get('KAFKA_TOPIC'), consumer_group='me', callback=process_message)  


# this is the method which is called for each message, so put your logic here!
def process_message(message: KMEMessage):
    print(f"Processing message {message}", flush=True)
    print(f"Message: {message.message}", flush=True)
    print(f"Topic: {message.topic}", flush=True)
    print(f"Completion Topic: {message.completion_topic}", flush=True)
    # if you want to pass a "completion message" you populate it and return it
    # if you don't want to pass a "completion message" just return `KMEMessage(topic='')`
    return_message = KMEMessage(topic=message.completion_topic)
    return_message.message = "foo the bar"
    return return_message

if __name__ == '__main__':
    print("Starting", flush=True)
    kme_producer() # produce a message
    kme_consumer() # now consume said message
    print("Finished", flush=True)

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

kme-0.1.1.tar.gz (4.4 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

kme-0.1.1-py3-none-any.whl (4.6 kB view details)

Uploaded Python 3

File details

Details for the file kme-0.1.1.tar.gz.

File metadata

  • Download URL: kme-0.1.1.tar.gz
  • Upload date:
  • Size: 4.4 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/3.4.1 importlib_metadata/4.5.0 pkginfo/1.7.0 requests/2.25.1 requests-toolbelt/0.9.1 tqdm/4.61.0 CPython/3.9.5

File hashes

Hashes for kme-0.1.1.tar.gz
Algorithm Hash digest
SHA256 8b82ea742dd05c89009a6e1a7f75fa4ce3efcb4cb81e3f32040bdfa72aed2674
MD5 078e75f5f51d84df97e6c174199e2063
BLAKE2b-256 cbeeb459be54e395e0347409acb38f4981fb3e7181d0fec904485917904c466d

See more details on using hashes here.

File details

Details for the file kme-0.1.1-py3-none-any.whl.

File metadata

  • Download URL: kme-0.1.1-py3-none-any.whl
  • Upload date:
  • Size: 4.6 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/3.4.1 importlib_metadata/4.5.0 pkginfo/1.7.0 requests/2.25.1 requests-toolbelt/0.9.1 tqdm/4.61.0 CPython/3.9.5

File hashes

Hashes for kme-0.1.1-py3-none-any.whl
Algorithm Hash digest
SHA256 63714067f75981567a1384a738093a05634928032b0d0ed19e21fd7a35a3b78c
MD5 cc7b93be611891671335f59fd7b27438
BLAKE2b-256 5eaf866ec438a66fda13e0cba147c7bc2cb41662517fc01f5b1064506abd0e5d

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page