单一位置的 Spark 模式管理

VB_*_*VB_ 5 hive apache-spark databricks aws-glue

管理 Spark 表模式的最佳方法是什么?您是否看到选项 2 的任何缺点?你能提出更好的选择吗?

我看到的解决方案

选项 1:为代码和元存储保留单独的定义

这种方法的缺点是您一直保持它们同步(容易出错)。另一个缺点 - 如果表有 500 列,它会变得很麻烦。

create_some_table.sql [第一个定义]

-- Databricks syntax (internal metastore)
CREATE TABLE IF NOT EXISTS some_table (
  Id int,
  Value string,
  ...
  Year int
)
USING PARQUET
PARTITION BY (Year)
OPTIONS (
  PATH 'abfss://...'
)
Run Code Online (Sandbox Code Playgroud)

some_job.py [第二定义]

def run():
   df = spark.read.table('input_table')  # 500 columns
   df = transorm(df)
   # this logic should be in `transform`, but anycase it should be
   df = df.select(
     'Id', 'Year', F.col('Value').cast(StringType()).alias('Value')  # actually another schema definition: you have to enumerate all output columns
   )
   df.write.saveAsTable('some_table')
Run Code Online (Sandbox Code Playgroud)

test_some_job.py [第三个定义]

def test_some_job(spark):
   output_schema = ...  # another definition
   expected = spark.createDataFrame([...], output_schema)
Run Code Online (Sandbox Code Playgroud)

选项 2:在代码中只保留一个定义(StructType)

可以动态生成模式。这种方法的好处 - 是简单和单一位置的模式定义。你看到任何缺点吗?

def run(input: Table, output: Table):
   df = spark.read.table(input.name)
   df = transform(df)
   save(df, output)    

def save(df: DataFrame, table: Table): 
    df \
        .select(table.schema.fieldNames()) \
        .write \
        .partitionBy(table.partition_by) \
        .option('path', table.path) \
        .saveAsTable(table.name)
    # In case table doesn't exists, Databricks will automatically generate table definition
        
class Table(NamedTuple):
    name: str
    path: str
    partition_by: List[str]
    schema: StructType
Run Code Online (Sandbox Code Playgroud)

Dou*_*s M 5

我先说几点,然后提出建议。

  1. 数据的寿命比代码长得多。
  2. 上面描述的代码是创建和写入数据的代码,还需要考虑读取和使用数据的代码。
  3. 还有第三个选项,将数据(模式)的定义与数据一起存储。通常称为“自描述格式”
  4. 数据的结构会随着时间的推移而改变。
  5. 这个问题被标记为databricksaws-glue
  6. Parquet 是在逐个文件的基础上进行自我描述的。
  7. Delta Lake 表使用 parquet 数据文件,但另外将架构嵌入到事务日志中,因此整个表和架构都是版本化的。
  8. 数据需要被广泛的工具生态系统使用,因此数据需要是可发现的,模式不应该被锁定在一个计算引擎中。

推荐:

  1. 以开放格式存储架构和数据
  2. 使用 Delta Lake 格式(结合了 Parquet 和事务日志)
  3. 改成USING PARQUETUSING DELTA
  4. 将您的元存储指向 AWS Glue Catalog,Glue Catalog 将存储表名称和位置
  5. 消费者将从 Delta Lake 表事务日志中解析架构
  6. 模式可以随着编写者代码的发展而发展。

结果:

  1. 您的作者创建模式,并且可以选择改进模式
  2. 所有消费者都将在 Delta Lake 中找到该架构(与表版本配对)(具体为 _delta_log 目录)