【问题标题】:Apply Rank or partitioned row_num function in Data Fusion在数据融合中应用 Rank 或分区 row_num 函数
【发布时间】:2020-09-24 00:52:50
【问题描述】:

我想在 Data Fusion 中对我的数据实现 rank 或分区 row_num 函数,但我没有找到任何插件。

有什么办法吗?

我想实现以下,

假设我有上面的数据,现在我想根据 AccountNumber 对数据进行分组,并将最近的记录发送到一个接收器并休息到其他接收器。 所以从上面的数据来看,

Sink1 预计有,

水槽2,

我计划通过按 AccountNumber 应用 rank 或 row_number 分区并按 Record_date desc 排序类似功能来进行这种隔离,并将 rank=1 或 row_num=1 的记录发送到一个接收器,然后发送到另一个接收器。

【问题讨论】:

  • 您在寻找任何具体的方法吗?
  • 嗨,Esteves,我无法在此处提供详细信息,因此我已用更多详细信息更新了问题。请检查

标签: google-cloud-data-fusion cdap


【解决方案1】:

解决问题的好方法是使用Spark plugin。 要将其添加到您的 Datafusion 实例,请转到 HUB -> 插件 -> 搜索 Spark -> 部署插件。然后您可以在 Analytics 选项卡上找到它。

为了给您举例说明如何使用它,我创建了以下管道:

这条管道基本上是:

  1. 从 GCS 读取文件。
  2. 在您的数据中执行排名函数
  3. 在不同的分支中过滤 rank=1 和 rank>1 的数据
  4. 将您的数据保存在不同的位置

现在让我们更深入地了解每个组件:

1 - GCS:这是一个简单的 GCS 源。本例使用的文件数据如下所示

2 - Spark_rank:这是一个 Spark 插件,代码如下。该代码基本上使用您的数据创建了一个临时视图,并且它们应用查询来对您的行进行排名。之后,您的数据将返回到管道。您还可以在下面看到此步骤的输入和输出数据。请注意输出是重复的,因为它被传递到两个分支。

      def transform(df: DataFrame, context: SparkExecutionPluginContext) : DataFrame = {
          df.createTempView("source")
          df.sparkSession.sql("SELECT AccountNumber, Address, Record_date, RANK() OVER (PARTITION BY accountNumber ORDER BY record_date DESC) as rank FROM source")
    }

3 - Spark2Spark3:与下面的步骤一样,此步骤使用 Spark 插件来转换数据。 Spark2 使用下面的代码只获取 rank = 1 的数据

    def transform(df: DataFrame, context: SparkExecutionPluginContext) : DataFrame = {
      df.createTempView("source_0")
      df.sparkSession.sql("SELECT AccountNumber, Address, Record_date FROM 
    source_0 WHERE rank = 1")
    }

Spark3 使用以下代码获取 rank > 1 的数据:

    def transform(df: DataFrame, context: SparkExecutionPluginContext) : DataFrame = {
      df.createTempView("source_1")
      df.sparkSession.sql("SELECT accountNumber, address, record_date FROM source_1 WHERE rank > 1")
    }

4 - GCS2GCS3:最后,在这一步中,您的数据会再次保存到 GCS。

【讨论】:

  • 嗨 Esteves,该解决方案看起来很有希望,但我收到错误:def transform(df: DataFrame, context: SparkExecutionPluginContext) : DataFrame = { df.createTempView("source") } \n\r这给了我错误:处理请求时遇到异常:org/apache/spark/api/java/function/Function.
  • 也许你可以用 createOrReplaceTempView 替换 createTempView
  • 嗨,Esteves,我已将 createTempView 替换为 createOrReplaceTempView,但在验证我的函数时仍然出现以下错误:处理请求时遇到异常:org/apache/spark/api/java/function/Function
  • @SUDHIRGARG 您是否取消了插件配置中的“在部署时编译”选项?
  • 是的,Esteves,我将它设置为 False,但仍然面临同样的问题。
猜你喜欢
  • 1970-01-01
  • 2016-03-03
  • 1970-01-01
  • 1970-01-01
  • 2023-04-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-01-25
相关资源
最近更新 更多