【发布时间】:2019-05-30 20:14:02
【问题描述】:
我希望一次将来自 Kafka 主题的所有聚合数据写入 parquet 文件(或者至少最终以一个 parquet 文件结束)。
我运行一个单独的生产者应用程序,该应用程序将 50 条消息放在主题上。 数据是在消费者应用程序中按时间(1 天)汇总的,因此我需要收集 1 天的所有数据并进行计数。这有效并且是这样完成的:
Dataset<Row> df = spark.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", BOOTSTRAP_SERVER)
.option("subscribe", "test")
.option("startingOffsets", "latest")
.option("group.id", "test")
.option("failOnDataLoss", false)
.option("key.deserializer", "org.apache.kafka.common.serialization.IntegerDeserializer")
.option("value.deserializer", "org.apache.kafka.common.serialization.StringSerializer")
.load()
// LEFT OUT CODE FOR READABILITY
.withWatermark("timestamp", "1 minutes")
.groupBy(
functions.window(new Column("timestamp"), "1 day", "1 day"),
new Column("container_nummer"))
.count();
然后将结果写入一个 parquet 文件,如下所示:
StreamingQuery query = df.writeStream()
.format("parquet")
.option("truncate", "false")
.option("checkpointLocation", "/tmp/kafka-logs")
.start("/Users/**/kafka-path");
query.awaitTermination();
如果我把它写到控制台,我最终会在第 1 批中得到正确的每一天的计数。当尝试将它写入 parquet 时,我只会得到多个空的 parquet 文件。我是这样读的:
SparkSession spark = SparkSession
.builder()
.appName("test")
.config("spark.master", "local")
.config("spark.sql.session.timeZone", "UTC")
.getOrCreate();
Dataset<Row> df = spark.read()
.parquet("/Users/**/kafka-path/part-00000-dd416263-8db1-4166-b243-caba470adac7-c000.snappy.parquet");
df.explain();
df.show(20);
所有 parquet 文件似乎都是空的(与将它们写入控制台相反),上面的代码输出如下:
+------+----------------+-----+
|window|container_nummer|count|
+------+----------------+-----+
+------+----------------+-----+
我有两个问题:
- 我的 parquet 文件为空的可能原因是什么?
- 最后是否有可能拥有 1 个完整的镶木地板文件,其中包含所有数据?我想使用这些数据在不同的程序中提供机器学习模型。
注意:它不需要在生产环境中运行。我只是希望有人知道这个工作的方法..
提前致谢!
【问题讨论】:
-
您能否将聚合数据写入 Parquet 文件?
标签: apache-spark spark-structured-streaming