You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何跨模块设置的变量值无法传递到其他函数?

原因分析

这是因为多进程的内存空间完全独立。主进程中修改的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 22:05:26