【问题标题】:Enforcing Schema for PySpark JobPySpark 作业的实施模式
【发布时间】:2020-07-22 17:59:04
【问题描述】:

我们有一堆不同的 pyspark 作业,并且有一个用例,其中两个作业可以链接在一起形成单独的第三个作业。因此,如果您有一份工作“A”和另一份工作“B”,则这些工作可以链接在一起形成另一份工作“C”。显然,只有当“A”的输出数据帧模式与“B”的数据帧模式兼容时,才会发生这种情况。这正是数据集 api 所实现的,但不幸的是 spark 没有为 python 提供数据集,这是可以理解的,因为 python 是动态类型的。然而,随着类型的出现,我们可以实现某种编译时安全性,我可以想到几种方法来实现这一点,一种是通过组合,另一种是通过继承。
组合

from typing import TypeVar, Generic
from pyspark.sql import DataFrame
from enum import Enum

#Marker Enum
class Schema(Enum):
    pass

T = TypeVar("T", bound = Schema)

class CustomDataSet(Generic[T]):
    def __init__(self, dataframe: DataFrame) -> None:
        self.dataframe = dataframe

现在我可以使用 mypy 并使用 CustomDataSet[MySchema](dataframe) 之类的组合。这里的问题是 MySchema 不是数据框对象,这可能会让使用它的人感到困惑。

继承

import abc
from enum import Enum

#Marker Enum
class Schema(Enum):
    pass

class MyInterface(metaclass=abc.ABCMeta):
    
    @abc.abstractmethod
    def get_input_schema() -> Schema:
        raise NotImplementedError
    
    @abc.abstractmethod
    def get_output_schema() -> Schema:
        raise NotImplementedError

现在'A'和'B'这两个作业可以实现上述接口了。这似乎是一种相当冗长的做事方式,感觉不是继承的粉丝。

我的问题是,有没有一种更 Pythonic 的方式可以更好地使用组合?

【问题讨论】:

    标签: python apache-spark pyspark enums


    【解决方案1】:

    quinn 库中有一个 validate_schema() 方法,如果 DataFrame 架构与所需的不同,则会引发异常:

    下面是实际的函数:

    data = [("jose", 1), ("li", 2), ("luisa", 3)]
    source_df = spark.createDataFrame(data, ["name", "age"])
    required_schema = StructType([
        StructField("name", StringType(), True),
        StructField("city", StringType(), True),
    ])
    quinn.validate_schema(source_df, required_schema) # throws a DataFrameMissingStructFieldError
    

    您可以在运行每个作业后验证架构。您概述的其他方法似乎有点过度设计 - 认为最好尽可能坚持简单的功能。好问题。

    【讨论】:

      猜你喜欢
      • 2021-06-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-10-05
      • 2021-02-02
      • 1970-01-01
      相关资源
      最近更新 更多