【问题标题】:Optimizing & batching a Parquet/JDBC join优化和批处理 Parquet/JDBC 连接
【发布时间】:2019-03-10 07:18:17
【问题描述】:

我正在执行从 S3 parquet 数据到 JDBC (Postgres) 表的连接操作,使用 parquet 数据中的列到 JDBC 表的主键。我需要 JDBC 表中的一小部分(但仍然是一个很大的数量 - 总共数万或数十万行),然后我需要智能地对数据进行分区以供在执行程序中使用。

我对整个数据工程尤其是 Spark 还是个新手,所以请原谅(并假设!)我的无知。我不太关心处理时间而不是内存使用;我必须使内存使用量符合 Amazon Glue 限制。

有什么好的方法可以做到这一点?

我现有的想法:

理论上,我可以构造如下 SQL 查询:

select * from t1 where id = key1 UNION
select * from t1 where id = key2 UNION...

但是,这似乎很愚蠢。这个问题:Selecting multiple rows by ID, is there a faster way than WHERE IN 让我想到了将要拉到临时表的键写入临时表,将其与原始表连接并拉出结果;这似乎是执行上述操作的“正确”方式。但是,这似乎也是一个很常见的问题,我还没有找到现成的解决方案。

也有可能在最小/最大 UUID 值之间拉动,但问题是我拉动了多少额外的行,并且由于 UUID 是,AFAIK,随机分布在可能的 UUID 值中,我希望获得很多额外的行(在连接期间将被遗漏的行)。不过,这可能是一种对 JDBC 数据进行分区的有用方法。

我还不清楚 JDBC 数据是如何到达执行程序的;它可能(全部)通过驱动程序进程。

因此,尝试将其形式化为问题:

  1. 是否有适用于这种用法的现有配方?
  2. 我应该考虑 Spark 的哪些功能来实现此目的?
  3. 来自 JDBC 连接的数据的实际 Spark 数据流是什么?

【问题讨论】:

  • 开箱即用的 Spark 不提供类似的功能。但是,如果您通过Russell Spitzer 搜索外部帖子,在类似情况下如何实现有用的下推。

标签: apache-spark jdbc parquet aws-glue


【解决方案1】:

似乎最好的方法(到目前为止)是将要获取的行 ID 写入数据库上的临时表,与主表进行连接,然后读出结果(如链接答案中所述)。

理论上,这在 Spark 中是完全可行的;像

// PSUEDOCODE!
df.select("row_id").write.jdbc(<target db>, "ids_to_fetch")
databaseConnection.execute("create table output from (select * from ids_to_fetch join target_table on row_id = id)")
df = df.join(
  spark.read.jdbc(<target db>, "output")
)

这可能是最有效的方法,因为 (AFAIK) 它会将 ID 的写入 连接表的读取都交给执行程序,而不是尝试执行驱动程序中的大部分内容。

但是,现在我无法将临时表写入目标数据库,因此我在驱动程序中生成了一系列 select where in 语句,然后提取这些语句的结果。

【讨论】:

    【解决方案2】:

    Spark 并非旨在对驱动程序执行任何高性能操作,最好避免使用它。

    对于您的情况,我建议先将数据从 S3 加载到某个 DF。保留此数据框,因为稍后将需要它。

    然后您可以使用 ma​​p(row->).distinct()

    的组合从 S3 为您的键解析唯一值

    然后对每个分区中具有合理数量的键的键进行分区,以便对 JDBC 进行单个查询。你也可以坚持上面的结果,执行count()操作,然后repartition()。例如,单个分区中的项目不超过 1000 个。

    然后使用 ma​​pPartitions,编写一个类似“SELECT * FROM table WHERE key in”的查询。

    然后使用 spark flatMap 需要执行实际的选择。我不知道使用数据框的自动方式,因此您可能需要直接使用 JDBC 来执行选择和映射数据。您可以在工作机器上初始化一个 spring 框架,并使用 spring 数据扩展轻松将数据从 DB 加载到某些实体列表。

    现在您已经有了 DatasSet,其中包含来自集群中 Postgres 的所需数据。您可以通过 toDF() 从中创建一个数据框。可能需要对此处的列进行一些额外的映射,或者在上一步将数据映射到 Row 类型。

    所以,现在您有 2 个必需的数据帧,一个是初始保存的来自 S3 的数据,另一个是来自 Postgres 的数据,您可以使用 Dataframe.join 以标准方式加入它们。 p>

    注意:当要重复使用时,不要忘记使用 .persist() 持久化数据集和框架。否则,它将每次都重复所有数据检索步骤。

    【讨论】:

      猜你喜欢
      • 2013-11-22
      • 1970-01-01
      • 1970-01-01
      • 2019-06-24
      • 1970-01-01
      • 2019-08-16
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多