【发布时间】: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