【问题标题】:How to optimize ETL data pipeline for fault tolerance when using Spark and Redshift?使用 Spark 和 Redshift 时如何优化 ETL 数据管道以实现容错?
【发布时间】:2021-07-04 20:02:03
【问题描述】:

我正在使用 PySpark 编写一个大批量作业,该作业对 200 个表进行 ETL 并加载到 Amazon Redshift 中。 这 200 个表是从一个输入数据源创建的。因此,只有当数据成功加载到 ALL 200 个表中时,批处理作业才会成功。批处理作业每天运行,同时将每个日期的数据附加到表中。

为了容错、可靠性和幂等性,我当前的工作流程如下:

  1. 使用临时表。使用CREATE TEMP TABLE LIKE <target_table> 创建临时 Redshift 表
  2. 将数据转换并加载到临时表中。
  3. 对其他 200 个表重复 1-2。
  4. 开始BEGIN事务。
  5. 将临时表数据复制到目标表中 使用INSERT INTO <taget_table> SELECT * FROM <staging_table>
  6. END交易
  7. DROP所有临时表。

这样我可以保证如果第 3 步失败(这种可能性更大),我不必担心从原始表中删除部分数据。相反,我将简单地重新运行整个批处理作业,因为临时表在 JDBC 断开后被丢弃。

虽然它解决了大部分问题,但它并不优雅、老套,而且会耗费额外的时间。我想如果 Spark 和/或 Redshift 提供标准工具来解决 ETL 世界中这个非常常见的问题。

谢谢

【问题讨论】:

    标签: apache-spark amazon-redshift spark-redshift


    【解决方案1】:

    COPY 命令可以在事务块中。你只需要:

    1. 开始
    2. 将数据复制到所有表中
    3. 提交(如果成功)

    Redshift 将为所有其他查看者维护表的先前版本,并且他们对表的视图在 COMMIT 之前不会改变。

    您布置的进程的好处是,在事务运行期间,其他进程无法获得表上的排他锁(ALTER TABLE 等)。您的插入将比 COPY 运行得快,因此事务的打开时间将更短。仅当其他进程在 ETL 运行的同时修改表时才会出现问题,这通常不是一个好主意。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2022-01-13
      • 2018-08-07
      • 1970-01-01
      • 1970-01-01
      • 2018-05-03
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多