【发布时间】:2020-07-22 20:21:08
【问题描述】:
我正在尝试编写一个 Spark Structured Streaming 作业,该作业从 Kafka 主题读取并通过 writeStream 操作写入单独的路径(在执行一些转换之后)。但是,当我运行以下代码时,只有第一个 writeStream 被执行,第二个被忽略。
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
.option("subscribe", "topic1") \
.load()
write_one = df.writeStream \
.foreachBatch(lambda x, y: transform_and_write_to_zone_one(x,y)) \
.start() \
.awaitTermination()
// transform df to df2
write_two = df2.writeStream \
.foreachBatch(lambda x, y: transform_and_write_to_zone_two(x,y)) \
.start() \
.awaitTermination()
我最初认为我的问题与这个post 有关,但是在将我的代码更改为以下内容后:
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
.option("subscribe", "topic1") \
.load()
write_one = df.writeStream \
.foreachBatch(lambda x, y: transform_and_write_to_zone_one(x,y)) \
.start()
// transform df to df2
write_two = df2.writeStream \
.foreachBatch(lambda x, y: transform_and_write_to_zone_two(x,y)) \
.start()
write_one.awaitTermination()
write_two.awaitTermination()
我收到以下错误:
org.apache.spark.sql.AnalysisException: Queries with streaming sources must be executed with writeStream.start();;
我不确定为什么start() 和awaitTermination() 之间的附加代码会导致上述错误(但我认为这可能是一个单独的问题,在answer 中引用到上面的同一篇文章)。在同一个作业中调用多个 writeStream 操作的正确方法是什么?最好在foreachBatch 调用的函数中同时写入这两个函数,还是有更好的方法来实现这一点?
【问题讨论】:
标签: apache-spark pyspark spark-structured-streaming