【发布时间】: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