【问题标题】:streaming aggregate not writing into sink流聚合未写入接收器
【发布时间】:2019-09-27 13:05:14
【问题描述】:

我必须处理一些每天收到的文件。该信息具有主键(日期、client_id、operation_id)。所以我创建了一个流,它只将新数据附加到增量表中:

operations\
        .repartition('date')\
        .writeStream\
        .outputMode('append')\
        .trigger(once=True)\
        .option("checkpointLocation", "/mnt/sandbox/operations/_chk")\
        .format('delta')\
        .partitionBy('date')\
        .start('/mnt/sandbox/operations')

这工作正常,但我需要总结按(日期,client_id)分组的这些信息,所以我创建了另一个从这个操作表到新表的流。所以我尝试将我的date 字段转换为时间戳,这样我就可以在编写聚合流时使用附加模式:

import pyspark.sql.functions as F

summarized= spark.readStream.format('delta').load('/mnt/sandbox/operations')
summarized= summarized.withColumn('timestamp_date',F.to_timestamp('date'))
summarized= summarized.withWatermark('timestamp_date','1 second').groupBy('client_id','date','timestamp_date').agg(<lot of aggs>)

summarized\
        .repartition('date')\
        .writeStream\
        .outputMode('append')\
        .option("checkpointLocation", "/mnt/sandbox/summarized/_chk")\
        .trigger(once=True)\
        .format('delta')\
        .partitionBy('date')\
        .start('/mnt/sandbox/summarized')

此代码运行,但它没有在接收器中写入任何内容。

为什么不将结果写入接收器?

【问题讨论】:

  • 数据的频率是多少?您的timestamp_date 列的每个1 second 将存在多少行?
  • 数据是每天的,因为timestamp_date 是从date 转换而来的,我认为所有行都在同一秒内
  • 这是不正确的。我认为您不太了解如何使用window 函数
  • 它不会写入接收器,因为summarized 是空的。您本质上是在说,groupby 每秒并执行聚合。而您可能想按一天或其他时间分组

标签: pyspark spark-structured-streaming azure-databricks delta-lake


【解决方案1】:

这里可能有两个问题。

格式错误的日期输入

我很确定问题出在 F.to_timestamp('date') 上,由于输入格式错误,它给出了 null

如果是这样,withWatermark('timestamp_date','1 second') 永远不会被“物化”并且不会触发任何输出。

您能否spark.read.format('delta').load('/mnt/sandbox/operations')(到read 而不是readStream)看看转换是否给出了正确的值?

spark.\
  read.\ 
  format('delta').\
  load('/mnt/sandbox/operations').\
  withColumn('timestamp_date',F.to_timestamp('date')).\
  show

所有行使用相同的时间戳

withWatermark('timestamp_date','1 second') 也有可能没有完成(因此“完成”了一个聚合),因为所有行都来自同一个时间戳,所以时间不会提前。

您应该有具有不同时间戳的行,以便每个 timestamp_date 的时间概念可以超过 '1 second' 延迟窗口。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多