【问题标题】:Snowflake is not deducting partitioned by column in Parquet雪花不扣除 Parquet 中的列分区
【发布时间】:2021-12-07 18:42:47
【问题描述】:

我对 Snowflake 的新功能 -Infer Schema 表功能有疑问。 INFER SCHEMA 函数在 parquet 文件上表现出色,并返回正确的数据类型。但是,当 parquet 文件被分区并存储在 S3 中时,INFER SCHEMA 无法像处理 pyspark 数据帧那样发挥作用。

在DataFrames中,分区文件夹名称和值作为最后一列读取;有没有办法在 Snowflake Infer 架构中实现相同的结果?

例子:

@GregPavlik - 输入采用结构化拼花格式。当 parquet 文件在没有分区的情况下存储在 S3 中时,架构是完美派生的。

示例:{ “AGMT_GID”:1714844883, “AGMT_TRANS_GID”:640481290, "DT_RECEIVED": "20 302", “LATEST_TRANSACTION_CODE”:“我” }

Snowflake 推断架构为我提供了 4 个列名及其数据类型。

但是,如果 parquet 文件存储在分区中 - 如上图所示。

在 - LATEST_TRANSACTION_CODE =I/ 文件夹下,我会将文件保存为

示例:{ “AGMT_GID”:1714844883, “AGMT_TRANS_GID”:640481290, “DT_RECEIVED”:“20 302” }

在这种情况下,snowflake infer Schema 只提供了三列;但是,在 Pyspark 数据框中读取同一个文件会打印所有四列。

我想知道 Snowflake 中是否有一种解决方法来读取分区的 parquet 文件。

【问题讨论】:

  • 你能展示一个简单的 JSON 输入和你想从模式推断中看到的输出吗?试图澄清最后一列点 - 似乎您想要的内容可从阶段读取的元数据中获得,但需要确认细节。
  • @GregPavlik - 我在原始问题中添加了更多详细信息。如果它澄清了您的疑问,请告诉我。

标签: snowflake-cloud-data-platform parquet


【解决方案1】:

在处理分区镶木地板文件时遇到了雪花问题。 这个问题不仅发生在 infer_schema 中,Following flow 不会将按列分区作为雪花中的列进行扣除:

  • 从镶木地板复制到表中
  • 从镶木地板合并到表中
  • 从镶木地板中选择
  • 镶木地板的 INFER_SCHEMA

Snowflake 将 parquet 文件视为文件并忽略文件夹名称中的元信息。 Apache Spark 智能扣除分区列。

以下方法是处理它的方法,直到雪花团队处理它。

方法 1

使用Snowflake metadata features 处理此问题。

截至目前,Snowflake 元数据仅提供

  • METADATA$FILENAME - 当前行所属的暂存数据文件的名称。包括阶段中数据文件的路径。
  • METADATA$FILE_ROW_NUMBER - 每条记录的行号

我们可以这样做:

    select $1:normal_column_1, ..., METADATA$FILENAME  
        FROM
            '@stage_name/path/to/data/' (pattern => '.*.parquet')
    limit 5;

这将给出一个包含分区文件完整路径的列。但是我们需要处理从中推导出列。例如:它会给出如下内容:

METADATA$FILENAME 
----------
path/to/data/year=2021/part-00020-6379b638-3f7e-461e-a77b-cfbcad6fc858.c000.snappy.parquet

我们可以做一个 regexp_replace 并将分区值作为这样的列:

    select 
        regexp_replace(METADATA$FILENAME, '.*\/year=(.*)\/.*', '\\1'
        ) as year
        $1:normal_column_1,  
    FROM
            '@stage_name/path/to/data/' (pattern => '.*.parquet')
    limit 5;
  • 在上面的正则表达式中,我们给出了分区键。
  • 第三个参数\\1是正则表达式组匹配号。在我们的例子中,第一个组匹配 - 这包含分区值。

方法 2

如果我们可以控制写入源 parquet 文件的流程。

  • 添加与按列内容分区具有相同内容的重复列
  • 这应该在写入 parquet 之前发生。所以parquet文件会有这个栏目内容。
df.withColumn("partition_column", col("col1")).write.partitionBy("partition_column").parquet(path)
  • 使用这种方法,如果我们这样做一次,在 parquet 的所有使用(COPY、MERGE、SELECT、INFER)中,新列将开始出现。

方法 3

如果我们无法控制写入源 parquet 文件的流程。

  • 这种方法更适合特定领域和数据模型。
  • 在许多用例中,我们需要对按列分区与数据的关系进行逆向工程。
  • 可以从其他列生成吗?比方说,如果数据是按年份分区的,其中年份是从 created_by 列派生的数据,那么这个派生数据可以再次重新生成。
  • 可以通过加入另一个雪花表来生成吗?假设 parquet 有一个 id,它可以与另一个表连接以在我们的列中动态派生

方法 3 更针对特定问题/领域。我们还需要在 parquet 的所有用例(COPY、MERGE、SELECT 等)中处理这个问题。

【讨论】:

    猜你喜欢
    • 2020-04-06
    • 1970-01-01
    • 1970-01-01
    • 2022-07-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-03
    相关资源
    最近更新 更多