openrvdas.logger.transforms.xml_aggregator_transform
No module-level documentation available.
1#!/usr/bin/env python3 2 3import logging 4 5from threading import Lock 6from xml.sax.handler import ContentHandler 7from xml.sax import make_parser 8 9from logger.transforms.transform import Transform # noqa: E402 10 11 12################################################################################ 13class XMLAggregatorTransform(Transform): 14 """Aggregate passed lines of XML until a complete XML record whose 15 outermost element matches 'tag' has been seen, then pass it on as a 16 single record.""" 17 ############################ 18 19 def __init__(self, tag, **kwargs): 20 """ 21 'tag' should be the identity of the top-level XML element that 22 we're expecting to read, e.g. 'OSU_DAS_Record'. 23 """ 24 super().__init__(**kwargs) # processes 'quiet' and type hints 25 self.tag = tag 26 27 # Only let one thread touch buffer at a time. Of course, if we're 28 # getting interleaved lines from different XML records here, we're 29 # screwed anyway. 30 self.buffer_lock = Lock() 31 self.buffer = '' 32 33 self.handler = _XMLHandler(tag=tag) 34 self.parser = make_parser(['xml.sax.IncrementalParser']) 35 self.parser.setContentHandler(self.handler) 36 37 ############################ 38 def transform(self, record: str): 39 """Aggregate, returning None until we're done, then return record.""" 40 41 # See if it's something we can process, and if not, try digesting 42 if not self.can_process_record(record): # inherited from BaseModule() 43 return self.digest_record(record) # inherited from BaseModule() 44 45 with self.buffer_lock: 46 # Feed record to the incremental parser 47 self.buffer += record + '\n' 48 self.parser.feed(record) 49 logging.debug('transform() got line: %s', record) 50 51 # If the record completes and XML record, it will be added to the 52 # queue in self.handler.items() - pop it off and return 53 if self.handler.complete(): 54 xml_record = self.buffer 55 self.buffer = '' 56 self.parser.close() 57 self.parser.reset() 58 self.handler.reset() 59 logging.debug('transform() got closing tag: %s', xml_record) 60 return xml_record 61 62 # Otherwise go home emptyhanded 63 return None 64 65 66################################################################################ 67class _XMLHandler(ContentHandler): 68 """ 69 Helper class here - _XMLHandler.endElement() will get called on 70 closing tags, so we can detect when we've got a closing tag for 71 whatever XML element we're after. We're omitting startElement, 72 and characters methods to store data on a stack during processing. 73 """ 74 ############################ 75 76 def __init__(self, tag): 77 super().__init__() 78 self.tag = tag 79 self.item_list = [] 80 self.element_complete = False 81 82 ############################ 83 def complete(self): 84 return self.element_complete 85 86 ############################ 87 def reset(self): 88 self.element_complete = False 89 90 ############################ 91 def endElement(self, name): 92 if name == self.tag: 93 # Create item from stored data on stack 94 self.element_complete = True
class
XMLAggregatorTransform(logger.transforms.transform.Transform):
14class XMLAggregatorTransform(Transform): 15 """Aggregate passed lines of XML until a complete XML record whose 16 outermost element matches 'tag' has been seen, then pass it on as a 17 single record.""" 18 ############################ 19 20 def __init__(self, tag, **kwargs): 21 """ 22 'tag' should be the identity of the top-level XML element that 23 we're expecting to read, e.g. 'OSU_DAS_Record'. 24 """ 25 super().__init__(**kwargs) # processes 'quiet' and type hints 26 self.tag = tag 27 28 # Only let one thread touch buffer at a time. Of course, if we're 29 # getting interleaved lines from different XML records here, we're 30 # screwed anyway. 31 self.buffer_lock = Lock() 32 self.buffer = '' 33 34 self.handler = _XMLHandler(tag=tag) 35 self.parser = make_parser(['xml.sax.IncrementalParser']) 36 self.parser.setContentHandler(self.handler) 37 38 ############################ 39 def transform(self, record: str): 40 """Aggregate, returning None until we're done, then return record.""" 41 42 # See if it's something we can process, and if not, try digesting 43 if not self.can_process_record(record): # inherited from BaseModule() 44 return self.digest_record(record) # inherited from BaseModule() 45 46 with self.buffer_lock: 47 # Feed record to the incremental parser 48 self.buffer += record + '\n' 49 self.parser.feed(record) 50 logging.debug('transform() got line: %s', record) 51 52 # If the record completes and XML record, it will be added to the 53 # queue in self.handler.items() - pop it off and return 54 if self.handler.complete(): 55 xml_record = self.buffer 56 self.buffer = '' 57 self.parser.close() 58 self.parser.reset() 59 self.handler.reset() 60 logging.debug('transform() got closing tag: %s', xml_record) 61 return xml_record 62 63 # Otherwise go home emptyhanded 64 return None
Aggregate passed lines of XML until a complete XML record whose outermost element matches 'tag' has been seen, then pass it on as a single record.
XMLAggregatorTransform(tag, **kwargs)
20 def __init__(self, tag, **kwargs): 21 """ 22 'tag' should be the identity of the top-level XML element that 23 we're expecting to read, e.g. 'OSU_DAS_Record'. 24 """ 25 super().__init__(**kwargs) # processes 'quiet' and type hints 26 self.tag = tag 27 28 # Only let one thread touch buffer at a time. Of course, if we're 29 # getting interleaved lines from different XML records here, we're 30 # screwed anyway. 31 self.buffer_lock = Lock() 32 self.buffer = '' 33 34 self.handler = _XMLHandler(tag=tag) 35 self.parser = make_parser(['xml.sax.IncrementalParser']) 36 self.parser.setContentHandler(self.handler)
'tag' should be the identity of the top-level XML element that we're expecting to read, e.g. 'OSU_DAS_Record'.
def
transform(self, record: str):
39 def transform(self, record: str): 40 """Aggregate, returning None until we're done, then return record.""" 41 42 # See if it's something we can process, and if not, try digesting 43 if not self.can_process_record(record): # inherited from BaseModule() 44 return self.digest_record(record) # inherited from BaseModule() 45 46 with self.buffer_lock: 47 # Feed record to the incremental parser 48 self.buffer += record + '\n' 49 self.parser.feed(record) 50 logging.debug('transform() got line: %s', record) 51 52 # If the record completes and XML record, it will be added to the 53 # queue in self.handler.items() - pop it off and return 54 if self.handler.complete(): 55 xml_record = self.buffer 56 self.buffer = '' 57 self.parser.close() 58 self.parser.reset() 59 self.handler.reset() 60 logging.debug('transform() got closing tag: %s', xml_record) 61 return xml_record 62 63 # Otherwise go home emptyhanded 64 return None
Aggregate, returning None until we're done, then return record.