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