【问题标题】:Writing all data at once to parquet file using structured streaming使用结构化流一次将所有数据写入镶木地板文件
【发布时间】: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


【解决方案1】:

读取正在进行的流所涉及的组件是结构化流。因此,您基本上是在读取无限制的消息/记录流,这些消息/记录必须在数据到达时写入。以此为基础,Spark 将继续分配执行器来执行读/写操作,并且将创建多个文件作为此过程的一部分。因此,您不会有一个包含所有数据的文件。以下是可用于写入 parquet 文件的语法:

df..writeStream.queryName("Loantxns_view").outputMode("append").format("parquet").option("path", "/user/root/Ln_sink2").option("checkpointLocation", "/user/root/Checkpoints2").start()

尝试使用以下机制读取parquet文件,看看是否得到了你要找的数据(ParquetDF是通过读取parquet目录路径读取的新数据帧):

val ParquetDF = spark.read.parquet("/user/root/Ln_sink2")
    ParquetDF.createOrReplaceTempView("xxxx_view");
    spark.sql("select * from xxxx_view").show(false);

试试这个,看看数据是否可见。

【讨论】:

    猜你喜欢
    • 2019-01-20
    • 1970-01-01
    • 2021-08-27
    • 2020-04-02
    • 2019-06-02
    • 2021-10-28
    • 1970-01-01
    • 1970-01-01
    • 2020-02-25
    相关资源
    最近更新 更多