【发布时间】:2017-09-26 21:25:10
【问题描述】:
我正在使用 Spark 结构化流从 Kafka 队列中读取数据。从 Kafka 阅读后,我在 dataframe 上应用 filter。我将此过滤后的数据框保存到镶木地板文件中。这会生成许多空的镶木地板文件。有什么办法可以停止写一个空文件?
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", KafkaServer) \
.option("subscribe", KafkaTopics) \
.load()
Transaction_DF = df.selectExpr("CAST(value AS STRING)")
decompDF = Transaction_DF.select(zip_extract("value").alias("decompress"))
filterDF = decomDF.filter(.....)
query = filterDF .writeStream \
.option("path", outputpath) \
.option("checkpointLocation", RawXMLCheckpoint) \
.start()
【问题讨论】:
-
这与我面临的类似挑战然后我决定将数据写入 hbase 然后向下流使用 hbase。在您的情况下,您也可以将其写入一些 Nosql 数据库以避免许多小文件。
标签: apache-spark pyspark spark-structured-streaming