5a0667beaf
git-subtree-dir: utils_v2 git-subtree-split: e7073350338ed6f28a4fdfab2e48c0bf458561b5
555 lines
20 KiB
Python
555 lines
20 KiB
Python
"""
|
|
|
|
AUTHOR:
|
|
|
|
Khushal P Soonderji
|
|
|
|
DATE:
|
|
|
|
Saturday, 13th Jul, 2024
|
|
|
|
OBJECTIVE:
|
|
|
|
To be able to work with keys
|
|
|
|
REFERENCES:
|
|
|
|
N/A
|
|
|
|
DOWNLOADS:
|
|
|
|
N/A
|
|
|
|
"""
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** IMPORT ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
# To make sibling directories accessible for imports:
|
|
import sys
|
|
sys.path.append(".")
|
|
sys.path.append("..")
|
|
|
|
# To use Kafka:
|
|
from aiokafka import AIOKafkaProducer
|
|
from aiokafka import AIOKafkaConsumer
|
|
|
|
# For working with JSON strings:
|
|
from utils_v2.string import json
|
|
from utils_v2.serialization.json_serializer import JSONSerializer
|
|
|
|
# For debugging:
|
|
from icecream import IceCreamDebugger
|
|
|
|
# For SSL security:
|
|
import ssl
|
|
|
|
# For asynchronous activities:
|
|
import asyncio
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** MACROS / ONE-TIME INIT ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
# --- Nothing Yet
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** VARIABLES ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
# --- Nothing Yet
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** FUNCTIONS ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
def get_ssl_context(
|
|
ca_file,
|
|
cert_file,
|
|
key_file
|
|
):
|
|
|
|
"""
|
|
Generate the SSL context to use with the Kafka instances.
|
|
:param ca_file: The Certificate Authority file as a path to a local file.
|
|
:param cert_file: The Certificate file as a path to a local file.
|
|
:param key_file: The Key file as a path to a local file.
|
|
:return: The SSL context instance as a path to a local file.
|
|
"""
|
|
|
|
ssl_context = ssl.create_default_context()
|
|
ssl_context.load_verify_locations(ca_file)
|
|
ssl_context.load_cert_chain(certfile = cert_file, keyfile = key_file)
|
|
return ssl_context
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** CLASSES ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
class ProducerKafka:
|
|
|
|
def __init__(
|
|
self,
|
|
topic,
|
|
serializer = None,
|
|
debug = True,
|
|
debug_prefix = "Kafka (P) | ",
|
|
**kwargs
|
|
):
|
|
|
|
"""
|
|
Create a Kafka Producer.
|
|
:param topic: The topic to produce on.
|
|
:param serializer: The serializer to use.
|
|
:param debug: Whether, or not, you want to print the debug strings.
|
|
:param debug_prefix: The prefix to use while debugging.
|
|
:param kwargs: Any configuration parameters for the Kafka instances.
|
|
"""
|
|
|
|
# Initialize the debugger:
|
|
self.__printer = IceCreamDebugger(prefix = debug_prefix, includeContext = True)
|
|
if not debug: self.__printer.disable()
|
|
|
|
# initialize the Kafka producer:
|
|
self.__topic = topic
|
|
self.__kwargs = kwargs
|
|
self.__producer = None
|
|
self.__connected = False
|
|
self.__serializer = serializer or JSONSerializer()
|
|
|
|
# For establishing connection:
|
|
self.__exclusive_semaphore = asyncio.Semaphore(1)
|
|
|
|
def enable_debug(self):
|
|
self.__printer.enable()
|
|
|
|
def disable_debug(self):
|
|
self.__printer.disable()
|
|
|
|
async def connect(self):
|
|
|
|
"""
|
|
Connects to the Kafka server if not connected.
|
|
:return: True or False based on the success of the operation.
|
|
"""
|
|
|
|
async with self.__exclusive_semaphore:
|
|
if not self.__connected:
|
|
try:
|
|
self.__producer = AIOKafkaProducer(**self.__kwargs)
|
|
await self.__producer.start()
|
|
self.__connected = True
|
|
except Exception as exception: self.__printer(exception)
|
|
return self.__connected
|
|
|
|
async def ensure_connection(self):
|
|
|
|
"""
|
|
Connects to the Kafka server if not connected.
|
|
:return: True or False based on the success of the operation.
|
|
"""
|
|
|
|
if not self.__connected: await self.connect()
|
|
return self.__connected
|
|
|
|
async def close(self):
|
|
|
|
"""
|
|
Terminates the connection.
|
|
:return: None.
|
|
"""
|
|
|
|
if self.__connected:
|
|
try:
|
|
await self.__producer.stop()
|
|
self.__printer("Producer closed!")
|
|
self.__connected = False
|
|
except Exception as exception: self.__printer(exception)
|
|
|
|
async def produce(self, value, key = None, topic = None, encoding = "utf-8"):
|
|
|
|
"""
|
|
Sends one message to the Kafka server on the topic that has been set for this instance.
|
|
:param value: The message to send.
|
|
:param key: The key to use when you want the messages to follow an order.
|
|
:param topic: A custom topic for this message, else the topic defined during the creation of this instance will
|
|
be used by default.
|
|
:param encoding: The encoding format.
|
|
:return: True or False based on the success of the operation.
|
|
"""
|
|
|
|
# Ensure connectivity to the server.
|
|
# If not connected, return with failure immediately.
|
|
if not await self.ensure_connection(): return False
|
|
|
|
try:
|
|
|
|
# Send the message:
|
|
await self.__producer.send_and_wait(
|
|
topic = topic or self.__topic,
|
|
value = self.__serializer.serialize(data = value, encoding = encoding),
|
|
key = key
|
|
)
|
|
|
|
# Return with success if no exception occurred:
|
|
return True
|
|
|
|
# Return with failure if something went wrong:
|
|
except Exception as exception:
|
|
self.__printer(exception, self.__topic, type(value), value)
|
|
return False
|
|
|
|
|
|
# ---------------------------------------------------------------------------------------------------------------------
|
|
|
|
|
|
class ConsumerKafka:
|
|
|
|
def __init__(
|
|
self,
|
|
topic,
|
|
serializer = None,
|
|
debug = True,
|
|
debug_prefix = "Kafka (C) | ",
|
|
**kwargs
|
|
):
|
|
|
|
"""
|
|
Create a Kafka Consumer.
|
|
:param topic: The topic to consumer on.
|
|
:param serializer: The serializer to use.
|
|
:param debug: Whether, or not, you want to print the debug strings.
|
|
:param debug_prefix: The prefix to use while debugging.
|
|
:param kwargs: Any configuration parameters for the Kafka instances.
|
|
"""
|
|
|
|
# Initialize the debugger:
|
|
self.__printer = IceCreamDebugger(prefix = debug_prefix, includeContext = True)
|
|
if not debug: self.__printer.disable()
|
|
|
|
# initialize the Kafka producer:
|
|
self.__topic = topic
|
|
self.__kwargs = kwargs
|
|
self.__consumer = None
|
|
self.__connected = False
|
|
self.__serializer = serializer or JSONSerializer()
|
|
|
|
# For establishing connection:
|
|
self.__exclusive_semaphore = asyncio.Semaphore(1)
|
|
|
|
def enable_debug(self):
|
|
self.__printer.enable()
|
|
|
|
def disable_debug(self):
|
|
self.__printer.disable()
|
|
|
|
async def connect(self):
|
|
|
|
"""
|
|
Connects to the Kafka server if not connected.
|
|
:return: True or False based on the success of the operation.
|
|
"""
|
|
|
|
async with self.__exclusive_semaphore:
|
|
if not self.__connected:
|
|
try:
|
|
self.__consumer = AIOKafkaConsumer(self.__topic, **self.__kwargs)
|
|
await self.__consumer.start()
|
|
self.__connected = True
|
|
except Exception as exception: self.__printer(exception)
|
|
return self.__connected
|
|
|
|
async def ensure_connection(self):
|
|
|
|
"""
|
|
Connects to the Kafka server if not connected.
|
|
:return: True or False based on the success of the operation.
|
|
"""
|
|
|
|
if not self.__connected: await self.connect()
|
|
return self.__connected
|
|
|
|
async def close(self):
|
|
|
|
"""
|
|
Terminates the connection.
|
|
:return: None.
|
|
"""
|
|
|
|
if self.__connected:
|
|
try:
|
|
await self.__consumer.stop()
|
|
self.__printer("Consumer closed!")
|
|
self.__connected = False
|
|
except Exception as exception: self.__printer(exception)
|
|
|
|
async def consume(self, count = 1, timeout = 0.05, encoding = "utf-8"):
|
|
|
|
"""
|
|
Get messages from the Kafka server.
|
|
:param count: The number of messages to get from the Kafka server.
|
|
:param timeout: The time in seconds to wait for retrieval.
|
|
:param encoding: The encoding to use.
|
|
:return: The messages that were received. If no messages are available, an empty list will be returned.
|
|
"""
|
|
|
|
# Ensure connectivity to the server.
|
|
# If not connected, return with failure immediately.
|
|
if not await self.ensure_connection(): return []
|
|
|
|
# Make a variable that will hold the final results:
|
|
messages = []
|
|
|
|
try:
|
|
|
|
# Read some messages:
|
|
results = await self.__consumer.getmany(
|
|
max_records = max(1, count),
|
|
timeout_ms = int(timeout * 1_000)
|
|
)
|
|
|
|
# Format the received messages:
|
|
if results:
|
|
for topic_partition, records in results.items():
|
|
for record in records:
|
|
record_dict = record.__dict__
|
|
record_dict["value"] = self.__serializer.deserialize(
|
|
data = record_dict["value"],
|
|
encoding = encoding
|
|
)
|
|
messages.append(record_dict)
|
|
|
|
# Debugging print if something went wrong:
|
|
except Exception as exception: self.__printer(exception)
|
|
|
|
# Done here:
|
|
return messages
|
|
|
|
|
|
# ---------------------------------------------------------------------------------------------------------------------
|
|
|
|
|
|
class BidirectionalKafka:
|
|
|
|
# The 'roles' that the instance can take.
|
|
# The master talks on the channel (topic) that the slave listens on and vice versa.
|
|
# Master-Slave is only for deciding who talks on which channel and who listens on which.
|
|
# In a two-party system, one must be the master, the other must be the slave.
|
|
# There are no extra privileges that the master enjoys. The naming convention was borrowed from common protocols
|
|
# used in electronics (like I2C).
|
|
ROLE_MASTER = 1
|
|
ROLE_SLAVE = 0
|
|
|
|
def __init__(
|
|
self,
|
|
role,
|
|
topic,
|
|
ack_topic: str = None,
|
|
group: str = None,
|
|
serializer = None,
|
|
debug = True,
|
|
debug_prefix = "Kafka (B) | ",
|
|
**kwargs
|
|
):
|
|
|
|
"""
|
|
Creates a walkie-talkie type setup to use Kafka in a bidirectional manner. Fo more information on all the
|
|
individual methods, please read through the doc-strings of the component classes 'ProducerKafka', and
|
|
'ConsumerKafka'.
|
|
:param role: Select from "ROLE_MASTER" and "ROLE_SLAVE". Between the two parties that are talking, one will be
|
|
the master and the other will be the slave. The channel that the master uses to speak will the one the slave
|
|
uses to listen, and vice versa.
|
|
:param topic: The topic to communicate on. Will be the same between the master and the slave.
|
|
:param ack_topic: Explicitly provide this for the second channel, or it will be created from the name of the
|
|
topic itself. Will be the same between the master and the slave.
|
|
:param group: The group to assign the instance to.
|
|
:param debug: Whether, or not, you want to print the debug strings.
|
|
:param debug_prefix: The prefix to use while debugging.
|
|
:param kwargs: Any configuration parameters for the Kafka instances.
|
|
"""
|
|
|
|
# Not down the basic variables:
|
|
self.__role = role
|
|
self.__topic = topic
|
|
self.__ack_topic = ack_topic or topic + "Ack"
|
|
self.__group = group
|
|
|
|
# In case the current instance is the master,
|
|
# it will talk on "topic", and listen on "ack_topic":
|
|
if self.__role == self.ROLE_MASTER:
|
|
self.__producer_kwargs = kwargs.copy()
|
|
self.__producer = ProducerKafka(
|
|
topic = self.__topic,
|
|
serializer = serializer,
|
|
debug = debug,
|
|
debug_prefix = debug_prefix.strip() + " (P) | ",
|
|
**self.__producer_kwargs
|
|
)
|
|
self.__consumer_kwargs = kwargs.copy()
|
|
self.__consumer_kwargs["group_id"] = self.__group
|
|
self.__consumer = ConsumerKafka(
|
|
topic = self.__ack_topic,
|
|
group = group,
|
|
serializer = serializer,
|
|
debug = debug,
|
|
debug_prefix = debug_prefix.strip() + " (C) | ",
|
|
**self.__consumer_kwargs
|
|
)
|
|
|
|
# On the other hand, if the current instance is a slave,
|
|
# It will listen on "topic", and talk on "ack_topic":
|
|
else:
|
|
self.__producer_kwargs = kwargs.copy()
|
|
self.__producer = ProducerKafka(
|
|
topic = self.__ack_topic,
|
|
serializer = serializer,
|
|
debug = debug,
|
|
debug_prefix = debug_prefix.strip() + " (P) | ",
|
|
**self.__producer_kwargs
|
|
)
|
|
self.__consumer_kwargs = kwargs.copy()
|
|
self.__consumer_kwargs["group_id"] = self.__group
|
|
self.__consumer = ConsumerKafka(
|
|
topic = self.__topic,
|
|
group = group,
|
|
serializer = serializer,
|
|
debug = debug,
|
|
debug_prefix = debug_prefix.strip() + " (C) | ",
|
|
**self.__consumer_kwargs
|
|
)
|
|
|
|
def enable_debug(self):
|
|
self.__producer.enable_debug()
|
|
self.__consumer.enable_debug()
|
|
|
|
def disable_debug(self):
|
|
self.__producer.disable_debug()
|
|
self.__consumer.disable_debug()
|
|
|
|
async def ensure_connection(self):
|
|
await self.__producer.ensure_connection()
|
|
await self.__consumer.ensure_connection()
|
|
|
|
async def close(self):
|
|
await self.__producer.close()
|
|
await self.__consumer.close()
|
|
|
|
async def produce(self, message, encoding = "utf-8"):
|
|
return await self.__producer.produce(message, encoding = encoding)
|
|
|
|
async def consume(self, count = 1, timeout = 0.05, encoding = "utf-8"):
|
|
return await self.__consumer.consume(count = count, timeout = timeout, encoding = encoding)
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** MAIN PROGRAM ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
if __name__ == "__main__":
|
|
|
|
import asyncio
|
|
import time
|
|
from data_models.kafka_message import KafkaMessage
|
|
|
|
ssl_ctx = get_ssl_context(
|
|
ca_file = r"/home/developer/PycharmProjects/utils/cred/kafka/cert_authority.pem",
|
|
cert_file = r"/home/developer/PycharmProjects/utils/cred/kafka/fullchain.pem",
|
|
key_file = r"/home/developer/PycharmProjects/utils/cred/kafka/privkey.pem"
|
|
)
|
|
|
|
async def consumer_test():
|
|
|
|
consumer = ConsumerKafka(
|
|
topic = "kft_file_upload",
|
|
# group_id = "assessImg",
|
|
group_id = "updateMedia",
|
|
bootstrap_servers = "wtt.ditscentre.in:9092",
|
|
security_protocol = "SSL",
|
|
ssl_context = ssl_ctx
|
|
)
|
|
await consumer.connect()
|
|
await asyncio.sleep(1.5)
|
|
print("READY!")
|
|
|
|
while True:
|
|
messages = await consumer.consume(count = 1)
|
|
if len(messages) > 0: print("MESSAGE:", json.to_string(messages[0], default=str))
|
|
await asyncio.sleep(1.0)
|
|
|
|
async def producer_test():
|
|
|
|
producer = ProducerKafka(
|
|
topic = "kft_file_upload",
|
|
bootstrap_servers = "del.ditscentre.in:9092",
|
|
security_protocol = "SSL",
|
|
ssl_context = ssl_ctx
|
|
)
|
|
await producer.connect()
|
|
print("READY!")
|
|
|
|
while True:
|
|
my_msg = KafkaMessage(
|
|
data = {
|
|
"accepted": False,
|
|
"reason": "low resolution"
|
|
},
|
|
media = {
|
|
"name": "pikachu_poster.jpg",
|
|
"ext": "jpg",
|
|
"url": "https://nexcom.ditscentre.in/utils/files/small/download/66ded1c1c1c05139a618b5ff",
|
|
"attr": {
|
|
"user": "SarangKabir",
|
|
"project": "ACE-PGP",
|
|
"id": 173,
|
|
"campaignActivityId": "25",
|
|
"idCampaign": 49,
|
|
"phoneNo": "7977821877"
|
|
}
|
|
},
|
|
appId = "aceWockhardt",
|
|
proc = {
|
|
"name": "_assessImg",
|
|
"attr": {
|
|
"blurThreshold": 0.25,
|
|
"clarityThreshold": 0.65,
|
|
"nsfwThreshold": 0.25,
|
|
"minWidth": 512,
|
|
"minHeight": 512
|
|
}
|
|
},
|
|
ack = None
|
|
)
|
|
success = await producer.produce(my_msg.model_dump())
|
|
print("produced...")
|
|
time.sleep(1.0)
|
|
break
|
|
|
|
await producer.close()
|
|
|
|
|
|
asyncio.run(consumer_test())
|