代码之家  ›  专栏  ›  技术社区  ›  RobinFrcd

使用PySpark UDF超时时返回None

  •  0
  • RobinFrcd  · 技术社区  · 6 年前

    我需要在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|
    +-----+-----+----+
    

    0 回复  |  直到 6 年前
    推荐文章