d5dc737b86
git-subtree-dir: utils_v2 git-subtree-split: 97ea8e2eee0b34c9bf906d60cd7fe14b4ef82daa
258 lines
8.8 KiB
Python
258 lines
8.8 KiB
Python
"""
|
|
|
|
AUTHOR:
|
|
|
|
Khushal P Soonderji
|
|
|
|
DATE:
|
|
|
|
Monday, 30th Dec., 2024.
|
|
|
|
OBJECTIVE:
|
|
|
|
To provide a standardized structure for Kafka messages.
|
|
|
|
REFERENCES:
|
|
|
|
N/A
|
|
|
|
DOWNLOADS:
|
|
|
|
N/A
|
|
|
|
"""
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** IMPORT ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
# To make sibling directories accessible for imports:
|
|
import sys
|
|
sys.path.append(".")
|
|
sys.path.append("..")
|
|
|
|
# For making data behaviour_models:
|
|
from pydantic import BaseModel, Field, field_validator, model_validator, AwareDatetime
|
|
from typing import Optional, Literal, Union, Dict, List, Any
|
|
|
|
# Related to Google:
|
|
from google.auth.transport.requests import Request
|
|
from google.oauth2.credentials import Credentials
|
|
|
|
# My utils:
|
|
from utils_v2.string import json
|
|
from utils_v2.string import regex
|
|
from utils_v2.date_time import date_time
|
|
|
|
# To work with date and time:
|
|
import datetime
|
|
import dateparser
|
|
|
|
# To make API calls:
|
|
import httpx
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** MACROS / ONE-TIME INIT ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
# --- Nothing Yet
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** VARIABLES ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
# --- Nothing Yet
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** FUNCTIONS ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
class ConsumedKafkaMessage(BaseModel):
|
|
|
|
topic: str = Field(
|
|
description = "the topic on which this message was received",
|
|
frozen = True
|
|
)
|
|
|
|
partition: int = Field(
|
|
description = "the partition in which this message was received",
|
|
frozen = True
|
|
)
|
|
|
|
offset: int = Field(
|
|
description = "the message's no. in the partition",
|
|
frozen = True
|
|
)
|
|
|
|
headers: List[Any] = Field(
|
|
description = "the headers received with the message",
|
|
frozen = True
|
|
)
|
|
|
|
key: Any = Field(
|
|
description = "the key with which this message is associated; important for partition management",
|
|
frozen = True
|
|
)
|
|
|
|
value: Any = Field(
|
|
description = "the actual payload of the message",
|
|
frozen = True
|
|
)
|
|
|
|
ts: AwareDatetime | None = Field(
|
|
description = "the time (utc) at which this message was sent to the queue",
|
|
frozen = True
|
|
)
|
|
|
|
tsType: Literal[
|
|
"createTime", # ...... The time at which the producer produced the message.
|
|
"logAppendTime", # ... The time at which the message was received by the broker.
|
|
None # ............... Unknown.
|
|
] = Field(
|
|
description = "to understand the source of the timestamp",
|
|
frozen = True
|
|
)
|
|
|
|
# ┏┓ ┏•
|
|
# ┃ ┏┓┏┓╋┓┏┓
|
|
# ┗┛┗┛┛┗┛┗┗┫
|
|
# ┛
|
|
|
|
class Config:
|
|
extra = "ignore"
|
|
populate_by_name = True
|
|
|
|
# ┓┏ ┓• ┓ •
|
|
# ┃┃┏┓┃┓┏┫┏┓╋┓┏┓┏┓
|
|
# ┗┛┗┻┗┗┗┻┗┻┗┗┗┛┛┗
|
|
|
|
@field_validator("ts", mode = "before")
|
|
def parse_dates(cls, value):
|
|
if value is None: return None
|
|
if not isinstance(value, datetime.datetime):
|
|
parsed = date_time.parse_date_time(
|
|
value,
|
|
date_formats = ["%Y-%m-%d %H:%M:%S"],
|
|
timezone = None
|
|
)
|
|
value = parsed if isinstance(parsed, datetime.datetime) else dateparser.parse(value)
|
|
if isinstance(value, datetime.datetime): value = date_time.to_timezone(value, date_time.TIMEZONE_UTC)
|
|
return value
|
|
|
|
# ┏┓ •
|
|
# ┃┃┏┓┏┓┏┓┏┓┏┓╋┓┏┓┏
|
|
# ┣┛┛ ┗┛┣┛┗ ┛ ┗┗┗ ┛
|
|
# ┛
|
|
|
|
pass
|
|
|
|
# ┏┓ ┏┓
|
|
# ┃ ┓┏┏╋┏┓┏┳┓ ┣ ┓┏┏┓┏┏
|
|
# ┗┛┗┻┛┗┗┛┛┗┗ ┻ ┗┻┛┗┗┛
|
|
|
|
@staticmethod
|
|
def from_aiokafka(
|
|
message,
|
|
deserializer = None
|
|
):
|
|
|
|
"""
|
|
To populate this model directly from the output of the 'aiokafka' library.
|
|
:param message: The raw message from the library.
|
|
:param deserializer: The function to use to deserialize the contents of the message.
|
|
:return: The standardized Kafka consumed message.
|
|
"""
|
|
|
|
# Extract the key and value:
|
|
key = message.key
|
|
value = message.value
|
|
if deserializer:
|
|
if key: key = deserializer(key)
|
|
if value: value = deserializer(value)
|
|
|
|
# Build and return the model:
|
|
return ConsumedKafkaMessage(
|
|
topic = message.topic,
|
|
partition = message.partition,
|
|
offset = message.offset,
|
|
headers = message.headers or [],
|
|
key = key,
|
|
value = value,
|
|
ts = message.timestamp / 1000.0 if message.timestamp else None,
|
|
tsType = {
|
|
0: "createTime",
|
|
1: "logAppendTime"
|
|
}.get(message.timestamp_type)
|
|
)
|
|
|
|
@staticmethod
|
|
def from_confluent_kafka(
|
|
message,
|
|
deserializer = None
|
|
):
|
|
|
|
"""
|
|
To populate this model directly from the output of the 'aiokafka' library.
|
|
:param message: The raw message from the library.
|
|
:param deserializer: The function to use to deserialize the contents of the message.
|
|
:return: The standardized Kafka consumed message.
|
|
"""
|
|
|
|
# If the message was null or an error:
|
|
if message is None or message.error(): return None
|
|
|
|
# Figure out the timestamp:
|
|
raw_ts = message.timestamp()
|
|
ts_type = raw_ts[0] if raw_ts else None
|
|
ts = raw_ts[1] / 1_000.0 if raw_ts else None
|
|
|
|
# Extract the key and value:
|
|
key = message.key()
|
|
value = message.value()
|
|
if deserializer:
|
|
if key: key = deserializer(key)
|
|
if value: value = deserializer(value)
|
|
|
|
# Build and return the model:
|
|
return ConsumedKafkaMessage(
|
|
topic = message.topic(),
|
|
partition = message.partition(),
|
|
offset = message.offset(),
|
|
headers = message.headers() or [],
|
|
key = key,
|
|
value = value,
|
|
ts = ts,
|
|
tsType = {
|
|
1: "createTime",
|
|
2: "logAppendTime"
|
|
}.get(ts_type)
|
|
)
|
|
|
|
|
|
# *****************************************************************************************************************
|
|
# ***** ****
|
|
# *** MAIN PROGRAM ***
|
|
# ***** ****
|
|
# *****************************************************************************************************************
|
|
|
|
|
|
if __name__ == "__main__":
|
|
|
|
pass
|