【问题标题】:Spark Streaming | Write different data frames to multiple tables in parallel火花流 |将不同的数据帧并行写入多个表
【发布时间】:2023-03-14 11:15:01
【问题描述】:

我正在从 Kafka 读取数据并加载到数据仓库中,我是从一个 Kafka 主题 创建一个数据框并在应用所需的转换后,我从中创建多个 DF 并将这些 DF 加载到不同的表中,但是此操作是按顺序进行的。有没有办法可以并行化这个表加载过程?

root
|-- attribute1Formatted: array (nullable = true)
|    |-- element: struct (containsNull = true)
|    |    |-- accexecattributes: struct (nullable = true)
|    |    |    |-- id: string (nullable = true)
|    |    |    |-- name: string (nullable = true)
|    |    |    |-- primary: boolean (nullable = true)
|    |    |-- accountExecUUID: string (nullable = true)
|-- attribute2Formatted: struct (nullable = true)
|    |-- Jake-DOT-Sandler@xyz.com: struct (nullable = true)
|    |    |-- id: string (nullable = true)
|    |    |-- name: string (nullable = true)
|    |    |-- primary: boolean (nullable = true)

分别为attribute1Formatted和attribute2Formatted创建了两个不同的数据框,并且这些DF被保存到不同表中的数据库中。

【问题讨论】:

  • 您可以发布您正在使用的代码吗?火花部分,简化输出业务部分

标签: scala dataframe apache-kafka spark-structured-streaming


【解决方案1】:

我对 Spark 流式传输了解不多,但我相信流式传输是迭代的微批处理,在 Spark 批处理执行中,每个操作都有一个接收器/输出。所以你不能一次执行就将它存储在不同的表中。

现在,

  1. 如果你把它写在一个表中,读者可以只读取他们需要的列。我的意思是:你真的需要把它存放在不同的地方吗?
  2. 可以写两次,过滤掉不需要的字段
  • 两个写入操作都会执行整个数据集的计算,然后删除不需要的列
  • 如果全数据集计算时间长,可以在过滤+写入前缓存

【讨论】:

  • 最初我正在写入一个表,但问题是 DF 有数组,并且在分解时会创建所有不同的行组合,我们希望将数组存储为单独的表。缓存我正在做,但我在想当表之间没有依赖关系时最好并行加载它们。
猜你喜欢
  • 2021-07-13
  • 1970-01-01
  • 2019-02-11
  • 2020-08-11
  • 1970-01-01
  • 2022-01-23
  • 2021-10-15
  • 2018-11-27
  • 2020-06-21
相关资源
最近更新 更多