codekingpro/portable-devtools
114k
1import base642import warnings3 4from scribe import scribe5from thrift.transport import TTransport, TSocket6from thrift.protocol import TBinaryProtocol7 8from eventlet import GreenPile9 10 11CATEGORY = 'zipkin'12 13 14class ZipkinClient(object):15 16 def __init__(self, host='127.0.0.1', port=9410):17 """18 :param host: zipkin collector IP address (default '127.0.0.1')19 :param port: zipkin collector port (default 9410)20 """21 self.host = host22 self.port = port23 self.pile = GreenPile(1)24 self._connect()25 26 def _connect(self):27 socket = TSocket.TSocket(self.host, self.port)28 self.transport = TTransport.TFramedTransport(socket)29 protocol = TBinaryProtocol.TBinaryProtocol(self.transport,30 False, False)31 self.scribe_client = scribe.Client(protocol)32 try:33 self.transport.open()34 except TTransport.TTransportException as e:35 warnings.warn(e.message)36 37 def _build_message(self, thrift_obj):38 trans = TTransport.TMemoryBuffer()39 protocol = TBinaryProtocol.TBinaryProtocolAccelerated(trans=trans)40 thrift_obj.write(protocol)41 return base64.b64encode(trans.getvalue())42 43 def send_to_collector(self, span):44 self.pile.spawn(self._send, span)45 46 def _send(self, span):47 log_entry = scribe.LogEntry(CATEGORY, self._build_message(span))48 try:49 self.scribe_client.Log([log_entry])50 except Exception as e:51 msg = 'ZipkinClient send error %s' % str(e)52 warnings.warn(msg)53 self._connect()54 55 def close(self):56 self.transport.close()57 