【问题标题】:Make another table with ID and distinct errors between two rows containing same id - SQL在包含相同 ID 的两行之间创建另一个 ID 和不同错误的表 - SQL
【发布时间】:2020-12-13 16:41:22
【问题描述】:

我有一个包含两列(“time_stamp”和“message”)的 spark 数据框。

示例数据框:

Time_stamp                   Message
2020-12-01 05:28:34:215      some text1 ID: 1
2020-12-01 05:28:40:210      some text2 error: A
2020-12-01 05:28:40:220      some text3 error: B
2020-12-01 05:28:41:203      some text4 error: A
2020-12-01 05:30:43:201      some text5 ID: 1
2020-12-01 05:32:50:215      some text6 ID: 2
2020-12-01 05:32:50:220      some text7 error: A
2020-12-01 05:48:51:220      some text8 error: C
2020-12-01 05:48:52:203      some text9 error: B
2020-12-01 05:51:53:201      some text10 ID: 2

我想在包含相同 ID 的两行之间创建另一个带有 ID 和明显错误的数据框。

预期输出:

示例表:

ID          Error
1           A
1           B
2           A
2           C
2           B

谢谢

【问题讨论】:

  • 我基本上是在数据块上使用 Apache spark 并拥有一个 spark 数据框。 SQL 查询可以在这个数据帧上运行。因此 mysql 查询可以轻松运行用于 spark 数据帧。使用 spark.sql
  • Spark SQL 与 MySQL 不同。您不能在 Spark 数据帧上运行 MySQL 查询。这些查询称为 Spark SQL 查询。
  • 哦,我明白了。感谢您的澄清。

标签: apache-spark pyspark apache-spark-sql


【解决方案1】:

试试下面的代码。代码按 ID 分组,收集错误消息并获取每个不同错误消息的最早错误消息。时间顺序保持不变。

import pyspark.sql.functions as F
from pyspark.sql.window import Window

df2 = df.withColumn(
    'Time_stamp',
    F.to_timestamp('Time_stamp', 'yyyy-MM-dd HH:mm:ss:SSS')
).withColumn(
    'ID',
    F.regexp_extract('Message', 'ID: ([a-zA-Z0-9]+)', 1)
).withColumn(
    'ID',
    F.last(F.when(F.col('ID') != '', F.col('ID')), True).over(Window.orderBy('Time_stamp'))
).filter(
    F.col('message').rlike('error')
).withColumn(
    'Message',
    F.regexp_extract('Message', 'error: (.*)', 1)
).groupBy('ID').agg(
    F.collect_set(F.array('Message', 'Time_stamp')).alias('Message')
).select(
    'ID',
    F.explode('Message').alias('Message')
).selectExpr(
    'ID',
    'Message[0] as error',
    'Message[1] as Time_stamp'
).withColumn(
    'rn',
    F.row_number().over(Window.partitionBy('ID', 'error').orderBy('Time_stamp'))
).filter('rn = 1').orderBy('Time_stamp').select('ID', 'error')

【讨论】:

  • 谢谢。让我理解代码并在实际数据中实现它。实际数据略有不同,需要相应修改。我应该更新我的数据以匹配实际数据吗?
  • @Dataholic 代码无需任何修改也应该可以工作。
  • 我运行了代码。输出数据框包含两列(ID 和错误),但这些列是空白的。
猜你喜欢
  • 2021-04-04
  • 1970-01-01
  • 2020-12-14
  • 1970-01-01
  • 2021-11-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-10-14
相关资源
最近更新 更多