【问题标题】:Spark Structured Streaming: join stream with data that should be read every micro batchSpark Structured Streaming:将流与应在每个微批次中读取的数据连接起来
【发布时间】:2018-07-26 03:33:25
【问题描述】:

我有一个来自 HDFS 的流,我需要将它与我也在 HDFS 中的元数据(两个 Parquet)加入。

我的元数据有时会更新,我需要加入最新的和最新的,这意味着理想情况下从 HDFS 读取每个流微批次的元数据。

我尝试对此进行测试,但不幸的是,即使我尝试使用 spark.sql.parquet.cacheMetadata=false,Spark 也会在缓存文件后读取元数据(据说)。

有没有办法读取每个微批次? Foreach Writer 不是我要找的?

以下是代码示例:

spark.sql("SET spark.sql.streaming.schemaInference=true")

spark.sql("SET spark.sql.parquet.cacheMetadata=false")

val stream = spark.readStream.parquet("/tmp/streaming/")

val metadata = spark.read.parquet("/tmp/metadata/")

val joinedStream = stream.join(metadata, Seq("id"))

joinedStream.writeStream.option("checkpointLocation", "/tmp/streaming-test/checkpoint").format("console").start()



/tmp/metadata/ got updated with spark append mode.

据我了解,通过 JDBC jdbc source and spark structured streaming 访问元数据,Spark 将查询每个微批次。

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-streaming


    【解决方案1】:

    据我所知,有两种选择:

    1. 创建临时视图并使用间隔刷新它:

      metadata.createOrReplaceTempView("metadata")

    并在单独的线程中触发刷新:

    spark.catalog.refreshTable("metadata")
    

    注意:在这种情况下,spark 将只读取相同的路径,如果您需要从 HDFS 上的不同文件夹读取元数据,则它不起作用,例如带有时间戳等。

    1. Tathagata Das suggested 的间隔重新启动流

    这种方式不适合我,因为我的元数据可能每小时刷新几次。

    【讨论】:

    猜你喜欢
    • 2021-07-13
    • 1970-01-01
    • 2021-04-26
    • 2020-07-22
    • 1970-01-01
    • 1970-01-01
    • 2019-06-25
    • 1970-01-01
    • 2019-10-26
    相关资源
    最近更新 更多