我需要在PySpark上运行长时间运行的任务(udf),有些任务可以运行几个小时,但我想添加一些超时包装器,以防它们运行得太长。我只想回去
None
如果有超时。
我做了点什么
signal
但我很确定这样做不是最安全的。
import pyspark
import signal
import time
from pyspark import SQLContext
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from pyspark.sql.functions import udf
conf = pyspark.SparkConf()
sc = pyspark.SparkContext.getOrCreate(conf=conf)
spark = SQLContext(sc)
schema = StructType([
StructField("sleep", IntegerType(), True),
StructField("value", StringType(), True),
])
data = [[1, "a"], [2, "b"], [3, "c"], [4, "d"], [1, "e"], [2, "f"]]
df = spark.createDataFrame(data, schema=schema)
def handler(signum, frame):
raise TimeoutError()
def squared_typed(s):
def run_timeout():
signal.signal(signal.SIGALRM, handler)
signal.alarm(3)
time.sleep(s)
return s * s
try:
return run_timeout()
except TimeoutError as e:
return None
squared_udf = udf(squared_typed, IntegerType())
df.withColumn('sq', squared_udf('sleep')).show()
它是有效的,给了我预期的输出,但是有没有一种方法可以在一个更大的
怎么了?
+-----+-----+----+
|sleep|value| sq|
+-----+-----+----+
| 1| a| 1|
| 2| b| 4|
| 3| c|null|
| 4| d|null|
| 1| e| 1|
| 2| f| 4|
+-----+-----+----+