【问题标题】:spark struct streaming writeStream output no data but no errorspark struct streaming writeStream 输出无数据但无错误
【发布时间】:2021-10-25 23:48:14
【问题描述】:

我有一个结构流作业,它从 Kafka 主题读取消息,然后保存到 dbfs。代码如下:

input_stream = spark.readStream \
  .format("kafka") \
  .options(**kafka_options) \
  .load() \
  .transform(create_raw_features)

# tranformation by 7 days rolling window
def transform_func(df):
  window_spec = window("event_timestamp", "7 days", "1 day")
  return df \
          .withWatermark(eventTime="event_timestamp", delayThreshold="2 days") \
          .groupBy(window_spec.alias("window"), "customer_id") \
          .agg(count("*").alias("count")) \
          .select("window.end", "customer_id", "count")

result = input_stream.transform(transform_func)

query = result \
    .writeStream \
    .format("memory") \
    .queryName("test") \
    .option("truncate","false").start()

我可以看到检查点工作正常。但是没有数据输出。

spark.table("test").show(truncate=False)

显示空表。有什么线索吗?

【问题讨论】:

  • 你等了 7 天吗?可能值得用较小的窗口大小测试代码。
  • spark.table("test") 在哪里运行?如果那是一个单独的 Spark 应用程序,那么我认为它不能从其他应用程序访问 .format("memory") 数据
  • @OneCricketeer 我在同一个 Databrick 笔记本中运行应用程序,因此它们共享同一个 SparkSession。

标签: pyspark apache-kafka spark-structured-streaming spark-kafka-integration spark3


【解决方案1】:

我发现了问题。在 Spark 文档output mode 部分中,它指出:

追加模式使用水印删除旧的聚合状态。但是窗口聚合的输出延迟了 withWatermark() 中指定的后期阈值,因为模式语义,行在最终确定后(即越过水印后)只能添加到结果表一次。

由于我没有明确指定输出模式,append 是隐式应用的,这意味着只有在水印阈值通过后才会出现第一个输出。

要获取每个微批次的输出,请改用输出模式updatecomplete

这对我有用

query = result \
    .writeStream \
    .format("memory") \
    .outputMode("update") \
    .queryName("test") \
    .option("truncate","false").start()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-08
    • 2017-12-20
    • 1970-01-01
    • 2017-04-04
    • 2017-03-28
    • 2016-09-10
    相关资源
    最近更新 更多