Skip to content

Commit

Permalink
add schema registry
Browse files Browse the repository at this point in the history
  • Loading branch information
sauljabin committed Jul 12, 2024
1 parent d7879f9 commit 994dde4
Show file tree
Hide file tree
Showing 2 changed files with 17 additions and 4 deletions.
1 change: 1 addition & 0 deletions kaskade/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,3 +34,4 @@ def get_kaskade_home() -> Path:
logger = logging.getLogger()
logger.addHandler(logger_handler)
logger.setLevel(logging.INFO)
logging.captureWarnings(True)
20 changes: 16 additions & 4 deletions kaskade/services.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,9 @@
ConfigSource,
)
from confluent_kafka.cimpl import NewTopic, NewPartitions
from confluent_kafka.schema_registry.protobuf import ProtobufDeserializer
from google.protobuf import descriptor_pb2, message_factory, json_format
from google.protobuf.message import DecodeError

from kaskade import logger
from kaskade.models import (
Expand Down Expand Up @@ -432,6 +434,8 @@ def sort_by_topic_name(topic: TopicMetadata) -> Any:
if __name__ == "__main__":

# protoc --include_imports --proto_path=. --python_out=. --descriptor_set_out=./Invoice.desc ./Invoice.proto
# https://docs.confluent.io/platform/current/schema-registry/fundamentals/serdes-develop/index.html#wire-format

file_desc = "/home/saul/Workspace/kafka-sandbox/kafka-protobuf/src/main/proto/Invoice.desc"
with open(file_desc, "rb") as f:
file_content = f.read()
Expand All @@ -442,9 +446,17 @@ def sort_by_topic_name(topic: TopicMetadata) -> Any:
def deserialize(raw_message: bytes | None) -> dict[str, Any] | None:
if raw_message is None:
return None
new_message = deserialization_class()
new_message.ParseFromString(raw_message)
return json_format.MessageToDict(new_message)

try:
new_message = deserialization_class()
new_message.ParseFromString(raw_message)
return json_format.MessageToDict(new_message, always_print_fields_with_no_presence=True)
except DecodeError:
protobuf_deserializer = ProtobufDeserializer(
deserialization_class, {"use.deprecated.format": False}
)
new_message = protobuf_deserializer(raw_message, None)
return json_format.MessageToDict(new_message, always_print_fields_with_no_presence=True)

async def main() -> None:
deserializer = DeserializerFactory()
Expand All @@ -454,7 +466,7 @@ async def main() -> None:
deserializer,
Format.STRING,
Format.BYTES,
page_size=5,
page_size=15,
)
messages = await consumer_service.consume()
for message in messages:
Expand Down

0 comments on commit 994dde4

Please sign in to comment.