import logging from concurrent.futures import ProcessPoolExecutor, as_completed from contextlib import closing from multiprocessing import Process, Queue from pickle import PicklingError, dumps, loads import pytest from structlog.threadlocal import bind_threadlocal from smelog.entities import LoggerConfig from smelog.factory import LoggerFactory def mock_func(logger, service_name, numb): logger = logger.bind(service_name=service_name) logger.info(numb) return True def test_multiprocessing_ok(): logger = LoggerFactory( LoggerConfig( name='test_config_mp_ok', version='0.0.1', level=logging.DEBUG, environment='dev', is_local=False, is_multiprocessing=True, ) ).get_logger('test_config_mp_ok') numbers = [1, 2, 3, 4, 5] bind_threadlocal(test=123) # all objects that are passed in submit() must be pickleable dumped = dumps(logger) logger = loads(dumped) try: with ProcessPoolExecutor(max_workers=4) as executor: futures = [ executor.submit(mock_func, logger, f'test_name{num}', num) for num in numbers ] for i, future in enumerate(as_completed(futures)): logger.info('iteration i: %s, progress: %s', i, future.result()) assert future.result() is True finally: logger.close() def test_multiprocessing_use_bind_taise_err(): logger = LoggerFactory( LoggerConfig( name='test_config_raise', version='0.0.1', level=logging.DEBUG, environment='dev', is_local=False, is_multiprocessing=True, ) ).get_logger('test_config_raise') logger = logger.bind(test=123) with closing(logger): with pytest.raises(PicklingError): dumps(logger) def writer(job_queue, result_queue, logger): try: while True: logger = logger.bind(test='123') value = job_queue.get() logger.info('%d', value) if value == -1: break result_queue.put(value * 2) except Exception as exc: # pylint: disable=broad-except logger.exception('shutdown: %s', exc) def test_multiprocessing_queue(): logger = LoggerFactory( LoggerConfig( name='test_config_queue', version='0.0.1', level=logging.DEBUG, environment='dev', is_local=False, is_multiprocessing=True, ) ).get_logger('test_config_queue') with closing(logger): job_queue = Queue() result_queue = Queue() proc = Process(target=writer, args=(job_queue, result_queue, logger)) proc.start() for i in range(10): job_queue.put(i) for i in range(10): result = result_queue.get() logging.info('result %d', result) job_queue.put(-1) proc.join()