-
Notifications
You must be signed in to change notification settings - Fork 2
/
message_converter.py
82 lines (62 loc) · 2.76 KB
/
message_converter.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
import json
import logging
from abc import ABC, abstractmethod
from fastavro import validation
from commons.schema import SchemaRetriever
logger = logging.getLogger('root')
class MessageConverter(ABC):
def __init__(self):
super().__init__()
@abstractmethod
def convert(self, msg, schema_name):
pass
@abstractmethod
def convert_all(self, msgs, schema_name):
pass
class AvroConverter(MessageConverter):
def __init__(self, schema_retriever: SchemaRetriever):
self.schema_retriever = schema_retriever
super().__init__()
def convert(self, msg, schema_name):
logger.debug(f'Converting {msg} using the class {self.__class__.__name__} and schema {schema_name}')
# TDOD Do the conversion.
return msg
def convert_all(self, msgs, schema_name):
logger.debug(f'Converting {msgs} using the class {self.__class__.__name__} and schema {schema_name}')
# TDOD Do the conversion.
return msgs
class AvroValidatedJsonConverter(MessageConverter):
def __init__(self, schema_retriever: SchemaRetriever):
self.schema_retriever = schema_retriever
super().__init__()
def convert(self, msg, schema_name):
logger.debug(f'Converting {msg} using the class {self.__class__.__name__} and schema {schema_name}')
if self.validate(msg, schema_name):
# For consistency we write a single message also as a list
return json.dumps([msg])
else:
logger.warning(
f'Validation failed for messages {msg}. Please look for any errors. Not publishing...')
return None
def convert_all(self, msgs, schema_name):
logger.debug(f'Converting {msgs} using the class {self.__class__.__name__} and schema {schema_name}')
if self.validate_all(msgs, schema_name):
return json.dumps(msgs)
else:
logger.warning(
f'Validation failed for messages {msgs}. Please look for any errors. Not publishing...')
return None
def validate(self, msg, schema_name):
logger.debug(f'Validating {msg} using the class {self.__class__.__name__} and schema {schema_name}')
schema = self.schema_retriever.get_schema(schema_name=schema_name)
if schema is None:
return False
else:
return validation.validate(msg, schema, raise_errors=True)
def validate_all(self, msgs, schema_name):
logger.debug(f'Validating {msgs} using the class {self.__class__.__name__} and schema {schema_name}')
schema = self.schema_retriever.get_schema(schema_name=schema_name)
if schema is None:
return False
else:
return validation.validate_many(msgs, schema, raise_errors=True)