为何跨模块设置的变量值无法传递到其他函数?
原因分析
这是因为多进程的内存空间完全独立。主进程中修改的settings.SCRIPT_ENV仅存在于主进程的内存中,multiprocessing.Pool创建的子进程会重新导入settings模块,此时模块内的SCRIPT_ENV还是初始值None。
解决方案
方案1:直接通过参数传递环境变量
避免依赖全局模块变量,将SCRIPT_ENV作为参数传递给子进程执行的函数:
import argparse import settings from multiprocessing import Pool import json import time from kafka import KafkaProducer def set_brokers_and_cert_path(script_env): brokers = None cert_full_path = None print(f"I am here man\n\n {script_env}") if script_env == "test": brokers = settings.TEST_BROKERS cert_full_path = settings.BASE_CERT_PATH + "test_env/" elif script_env == "accept": brokers = settings.ACCEPT_BROKERS cert_full_path = settings.BASE_CERT_PATH + "new_env/" return brokers, cert_full_path def func_b(args): partition, script_env = args kafka_brokers, cert_full_path = set_brokers_and_cert_path(script_env) producer = KafkaProducer(bootstrap_servers=kafka_brokers, security_protocol='SSL', ssl_check_hostname=True, ssl_cafile=cert_full_path +'cacert.pem', ssl_certfile=cert_full_path + 'certificate.pem', ssl_keyfile=cert_full_path + 'key.pem', max_block_ms=1200000, value_serializer=lambda v: json.dumps(v).encode('utf-8'), key_serializer=str.encode ) try: producer.send(settings.KAFKA_TOPIC, value="some val", key="some key", timestamp_ms=int(time.time()), headers=[], partition=partition) producer.flush() except AssertionError as e: print(f"Error in partition: {partition}, {e}") def main(script_env): with Pool(settings.NUM_PROCESSES) as p: # 将环境变量与分区号打包成元组传递 p.map(func_b, [(i, script_env) for i in range(0, 24)]) if __name__ == "__main__": parser = argparse.ArgumentParser() parser.add_argument("--env", help="Environment against which the script needs to run") args = parser.parse_args() if args.env: main(args.env) else: raise Exception("Please pass env argument. Ex: --env test/accept")
方案2:使用Pool的初始化函数统一设置子进程变量
利用Pool的initializer和initargs参数,在每个子进程启动时初始化settings变量:
import argparse import settings from multiprocessing import Pool import json import time from kafka import KafkaProducer def init_child_process(script_env): # 子进程启动时统一设置settings变量 settings.SCRIPT_ENV = script_env def set_brokers_and_cert_path(): brokers = None cert_full_path = None print(f"I am here man\n\n {settings.SCRIPT_ENV}") if settings.SCRIPT_ENV == "test": brokers = settings.TEST_BROKERS cert_full_path = settings.BASE_CERT_PATH + "test_env/" elif settings.SCRIPT_ENV == "accept": brokers = settings.ACCEPT_BROKERS cert_full_path = settings.BASE_CERT_PATH + "new_env/" return brokers, cert_full_path def func_b(partition): kafka_brokers, cert_full_path = set_brokers_and_cert_path() producer = KafkaProducer(bootstrap_servers=kafka_brokers, security_protocol='SSL', ssl_check_hostname=True, ssl_cafile=cert_full_path +'cacert.pem', ssl_certfile=cert_full_path + 'certificate.pem', ssl_keyfile=cert_full_path + 'key.pem', max_block_ms=1200000, value_serializer=lambda v: json.dumps(v).encode('utf-8'), key_serializer=str.encode ) try: producer.send(settings.KAFKA_TOPIC, value="some val", key="some key", timestamp_ms=int(time.time()), headers=[], partition=partition) producer.flush() except AssertionError as e: print(f"Error in partition: {partition}, {e}") def main(script_env): # 创建Pool时指定初始化函数和参数 with Pool(settings.NUM_PROCESSES, initializer=init_child_process, initargs=(script_env,)) as p: p.map(func_b, [i for i in range(0, 24)]) if __name__ == "__main__": parser = argparse.ArgumentParser() parser.add_argument("--env", help="Environment against which the script needs to run") args = parser.parse_args() if args.env: main(args.env) else: raise Exception("Please pass env argument. Ex: --env test/accept")
方案对比
- 方案1:逻辑直观,消除全局变量依赖,适合参数较少的场景;
- 方案2:适合需要在多个子进程函数中共享全局变量的场景,无需逐个传递参数。
内容的提问来源于stack exchange,提问作者Rajat Bhardwaj
相关产品推荐
相关产品推荐

