【问题标题】:Pyspark don't recognize env variables on the method passed as argument to foreach or foreachPartitionPyspark 无法识别作为参数传递给 foreach 或 foreachPartition 的方法上的环境变量
【发布时间】:2017-05-10 17:27:50
【问题描述】:

在下面的代码中,我尝试实例化 redis-py 在 URL 处使用 env 变量进行连接。问题是当我使用foreach or foreachPartition 时,#save_on_redis 方法中无法识别 env 变量。

我只是尝试在外面创建redis连接,但我收到“pickle.PicklingError: Can't pickle 'lock' object”,因为spark尝试同时运行这两个方法, 在所有节点上。

问题:如何在作为参数传递给 foreach 或 foreachPartition 的方法上使用环境变量?

import os
from pyspark.sql import SparkSession
import redis

spark = (SparkSession
        .builder
        .getOrCreate())

print "---------"
print os.getenv("REDIS_REPORTS_URL")
print "---------"

def save_on_redis(row):
    redis_ = redis.StrictRedis(host=os.getenv("REDIS_REPORTS_URL"), port=6379, db=0)
    print os.getenv("REDIS_REPORTS_URL")
    print redis_
    redis_.set("#teste#", "fagner")


df  = spark.createDataFrame([(0,1), (0,1), (0,2)], ["id", "score"])
df.foreach(save_on_redis)

【问题讨论】:

    标签: python apache-spark pyspark redis-py


    【解决方案1】:

    我建议您在驱动程序进程中获取 env 变量并将其作为 python 变量传递给工作进程,您可以在其中使用 os.putenv 设置您的环境

    例子:

    In [1]: import os
    
    In [2]: a = sc.parallelize(range(20))
    
    In [3]: os.getenv('MY_VAR')
    Out[3]: 'some_value'
    
    In [4]: def f(iter):
        import os
        return (str(os.getenv('MY_VAR')),)
       ...:
    
    In [5]: a.mapPartitions(f).collect()
    Out[5]: ['None', 'None']
    
    In [6]: my_var = os.getenv('MY_VAR')
    
    In [6]: def f2(iter):
        import os
        from subprocess import check_output
        os.putenv('MY_VAR', my_var)
        return (check_output('env | grep MY_VAR', shell=True), my_var)
       ....:
    
    In [7]: a.mapPartitions(f2).collect()
    Out[7]:
    ['MY_VAR=some_value\n',
     'some_value',
     'MY_VAR=some_value\n',
     'some_value']
    

    PS。根据this answer,最好直接修改os.environ映射对象,而不是使用os.putenv

    【讨论】:

      猜你喜欢
      • 2020-11-14
      • 1970-01-01
      • 2019-06-02
      • 2016-12-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多